#![allow(deprecated)]
pub mod application_crypto;
pub mod application_crypto_streams;
pub mod broadcast;
pub mod client;
pub mod coordination;
pub mod datagrams;
pub mod explicit_transfer_crypto;
pub(crate) mod generated;
pub mod heartbeat;
#[cfg(feature = "iroh-carrier-core")]
pub mod iroh_carrier;
#[cfg(feature = "iroh-carrier-core")]
pub mod iroh_carrier_bootstrap;
pub mod iroh_carrier_kind;
#[cfg(feature = "iroh-carrier-core")]
pub mod iroh_carrier_proof;
pub(crate) mod iroh_connection_policy;
pub mod key_agreement;
pub mod lifecycle_reason;
#[cfg(feature = "managed-group-encryption")]
#[allow(
dead_code,
reason = "gated managed-room controller; consumed only by managed gateway adapters"
)]
pub(crate) mod managed_group_controller;
#[cfg(feature = "managed-group-encryption")]
#[allow(
dead_code,
reason = "gated managed-room crypto owner; enabled only after the tracked persistence and router milestones"
)]
pub(crate) mod managed_group_crypto;
pub mod media;
pub mod native_protocol;
pub mod offline;
#[cfg(feature = "iroh-carrier-core")]
pub mod packet_carrier_transport;
pub mod presence;
pub(crate) mod presence_policy;
pub mod protocol_config;
pub mod route_policy;
pub mod runtime_policy;
pub mod service_errors;
pub mod session_token;
pub mod signaling;
pub mod sparse_fanout;
pub mod stream_metadata;
pub(crate) mod transport_generation;
pub(crate) mod transport_label;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
pub mod local_discovery;
#[cfg(test)]
pub mod test_constants {
pub const TEST_PROJECT_ID: &str = "test-project";
pub const TEST_API_KEY: &str = "pk_test_0000000000000000000000000000000000000000";
}
pub const LIVE_PROJECT_ID: &str = "pluto-rtc-prod";
pub fn validate_api_key(api_key: &str) -> anyhow::Result<&str> {
let trimmed = api_key.trim();
let valid_prefix = trimmed.starts_with("pk_live_") || trimmed.starts_with("pk_test_");
let suffix = trimmed.get(8..).unwrap_or_default();
if !valid_prefix || suffix.len() != 40 || !suffix.bytes().all(|byte| byte.is_ascii_hexdigit()) {
anyhow::bail!("OpenRTC 2.0 requires a public pk_live_ or pk_test_ API key");
}
Ok(trimmed)
}
pub fn app_tag_from_api_key(api_key: &str) -> String {
let trimmed = api_key.trim();
if trimmed.is_empty() {
return "app_anonymous".to_string();
}
let suffix_len = trimmed.len().min(16);
format!("app_{}", &trimmed[trimmed.len() - suffix_len..])
}
pub fn space_app_tag(api_key: &str, space_key: &str) -> String {
let input = format!("{}:{}", api_key.trim(), space_key.trim());
let digest = <sha2::Sha256 as sha2::Digest>::digest(input.as_bytes());
format!("space::{}", hex::encode(digest))
}
#[cfg(test)]
mod constructor_contract_tests {
use super::*;
#[test]
fn rust_constructor_is_provider_neutral_and_side_effect_free() {
let api_key = test_constants::TEST_API_KEY;
let client = client::Client::new(api_key.to_string()).expect("valid public API key");
assert_eq!(client.app_tag(), app_tag_from_api_key(api_key));
assert!(client::Client::new("firebase-project-id".to_string()).is_err());
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn ensure_rustls() {
if rustls::crypto::CryptoProvider::get_default().is_none() {
let _ = rustls::crypto::ring::default_provider().install_default();
}
}
#[cfg(not(target_arch = "wasm32"))]
pub mod adapters;
#[cfg(not(target_arch = "wasm32"))]
pub(crate) mod native_coordination_gateway;
#[cfg(not(target_arch = "wasm32"))]
pub mod native;
#[cfg(not(target_arch = "wasm32"))]
mod native_key_store;
pub use client::Client;
pub mod ticket_mesh;
#[cfg(not(target_arch = "wasm32"))]
pub use native::ControlPlane;
#[cfg(not(target_arch = "wasm32"))]
pub use ticket_mesh::TicketMesh;
pub use ticket_mesh::{TicketMeshError, TicketMeshErrorCode, TicketMeshOptions};
pub mod connection_manager;
pub mod protocol_registry;
#[cfg(not(target_arch = "wasm32"))]
pub mod runtime_manager;
#[cfg(not(target_arch = "wasm32"))]
pub mod transport;
#[cfg(not(target_arch = "wasm32"))]
pub use client::EndpointHandle;
#[cfg(all(
not(target_arch = "wasm32"),
not(any(target_os = "ios", target_os = "android"))
))]
pub mod sso;
#[cfg(not(target_arch = "wasm32"))]
pub mod native_node;
#[cfg(not(target_arch = "wasm32"))]
pub mod native_device;
#[cfg(all(
test,
not(target_arch = "wasm32"),
feature = "iroh-protocols-wasm",
any(feature = "transport-webrtc", feature = "transport-moq")
))]
mod native_carrier_protocol_test;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
pub mod native_moq_carrier;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
pub mod native_webrtc_carrier;
#[cfg(all(target_arch = "wasm32", feature = "iroh-protocols-wasm"))]
mod wasm_docs_persistence;
#[cfg(all(target_arch = "wasm32", feature = "iroh-protocols-wasm"))]
mod wasm_indexeddb_blob_store;
#[cfg(all(target_arch = "wasm32", feature = "transport-moq"))]
pub mod wasm_moq_carrier;
#[cfg(target_arch = "wasm32")]
pub mod wasm_node;
#[cfg(all(target_arch = "wasm32", feature = "transport-webrtc"))]
pub mod wasm_webrtc_carrier;
#[cfg(target_arch = "wasm32")]
#[macro_export]
macro_rules! console_log {
($($t:tt)*) => (web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format_args!($($t)*).to_string())))
}
#[cfg(target_arch = "wasm32")]
pub mod wasm_api {
use crate::client::Client;
use crate::session_token::split_ticket;
use crate::wasm_node::{peer_uni_stream_from_send, BiStream, PeerUniStream};
use iroh_tickets::endpoint::EndpointTicket;
use std::cell::RefCell;
use std::collections::{HashMap, HashSet};
use std::rc::Rc;
use std::str::FromStr;
use std::sync::{Arc, Mutex};
use wasm_bindgen::prelude::*;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
use wasm_bindgen_futures::spawn_local;
use wasm_streams::readable::sys::ReadableStream as JsReadableStream;
async fn send_native_signal_control_frame(
inner: &Client,
endpoint_id: String,
frame: Vec<u8>,
) -> Result<(), JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|error| JsValue::from_str(&format!("{error}")))?;
inner
.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let (mut send, _recv) = inner
.open_bi_internal(endpoint_id)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
send.write_all(&[0x00])
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
send.write_all(&(b"signal".len() as u32).to_be_bytes())
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
send.write_all(b"signal")
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
send.write_all(&frame)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
crate::application_crypto_streams::PeerSendStream::plain(send)
.finish_and_wait_for_peer(std::time::Duration::from_secs(2))
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "transport-webrtc")]
#[derive(Debug, Clone)]
struct WasmWebRtcCarrierAttempt {
connection_id: String,
remote_endpoint_id: iroh::EndpointId,
bootstrap: crate::iroh_carrier_bootstrap::CarrierBootstrapFrame,
generation: crate::client::WasmPeerDataGeneration,
role: &'static str,
prepared: bool,
retry_sent: bool,
offer_started: bool,
remote_ready: bool,
completion_started: bool,
retry_count: u8,
inbound_authorization_expires_at_ms: Option<f64>,
}
#[cfg(feature = "transport-moq")]
#[derive(Debug, Clone)]
struct WasmMoqCarrierAttempt {
connection_id: String,
remote_endpoint_id: iroh::EndpointId,
bootstrap: crate::iroh_carrier_bootstrap::CarrierBootstrapFrame,
generation: crate::client::WasmPeerDataGeneration,
role: &'static str,
prepared: bool,
retry_sent: bool,
retry_count: u8,
inbound_authorization_expires_at_ms: Option<f64>,
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct WasmRemoteCarrierCapabilities {
webrtc: bool,
moq: bool,
}
#[wasm_bindgen]
pub struct WasmPeerDatagramPolicy {
inner: RefCell<crate::datagrams::PeerDatagramPolicy>,
}
#[wasm_bindgen]
impl WasmPeerDatagramPolicy {
#[wasm_bindgen(js_name = setIncomingMaxAge)]
pub fn set_incoming_max_age(&self, value: Option<u32>, now_ms: f64) {
self.inner
.borrow_mut()
.set_incoming_max_age_ms(value.map(u64::from), now_ms.max(0.0) as u64);
}
#[wasm_bindgen(js_name = setOutgoingMaxAge)]
pub fn set_outgoing_max_age(&self, value: Option<u32>) {
self.inner
.borrow_mut()
.set_outgoing_max_age_ms(value.map(u64::from));
}
#[wasm_bindgen(js_name = setIncomingMaxBufferedDatagrams)]
pub fn set_incoming_max_buffered_datagrams(&self, value: u32) -> Result<(), JsValue> {
self.inner
.borrow_mut()
.set_incoming_max_buffered(value as usize)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = setOutgoingMaxBufferedDatagrams)]
pub fn set_outgoing_max_buffered_datagrams(&self, value: u32) -> Result<(), JsValue> {
self.inner
.borrow_mut()
.set_outgoing_max_buffered(value as usize)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = pushIncoming)]
pub fn push_incoming(&self, payload: Vec<u8>, now_ms: f64) -> Result<(), JsValue> {
self.inner
.borrow_mut()
.push_incoming(payload, now_ms.max(0.0) as u64)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = popIncoming)]
pub fn pop_incoming(&self, now_ms: f64) -> Option<Vec<u8>> {
self.inner.borrow_mut().pop_incoming(now_ms.max(0.0) as u64)
}
#[wasm_bindgen(js_name = pushOutgoing)]
pub fn push_outgoing(&self, payload: Vec<u8>, now_ms: f64) -> Result<bool, JsValue> {
self.inner
.borrow_mut()
.push_outgoing(payload, now_ms.max(0.0) as u64)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = popOutgoing)]
pub fn pop_outgoing(&self, now_ms: f64) -> Result<JsValue, JsValue> {
let Some(outgoing) = self.inner.borrow_mut().pop_outgoing(now_ms.max(0.0) as u64)
else {
return Ok(JsValue::NULL);
};
let result = js_sys::Object::new();
js_sys::Reflect::set(
&result,
&JsValue::from_str("payload"),
&js_sys::Uint8Array::from(outgoing.payload.as_slice()),
)?;
js_sys::Reflect::set(
&result,
&JsValue::from_str("remainingMaxAgeMs"),
&outgoing
.remaining_max_age_ms
.map(|value| JsValue::from_f64(value as f64))
.unwrap_or(JsValue::NULL),
)?;
Ok(result.into())
}
#[wasm_bindgen(js_name = recordSent)]
pub fn record_sent(&self, bytes: usize) {
self.inner.borrow_mut().record_sent(bytes);
}
#[wasm_bindgen(js_name = recordExpiredOutgoing)]
pub fn record_expired_outgoing(&self) {
self.inner.borrow_mut().record_expired_outgoing();
}
#[wasm_bindgen(js_name = recordSendFailure)]
pub fn record_send_failure(&self) {
self.inner.borrow_mut().record_send_failure();
}
#[wasm_bindgen(js_name = getStats)]
pub fn get_stats(&self, now_ms: f64) -> Result<JsValue, JsValue> {
let stats = self.inner.borrow_mut().stats(now_ms.max(0.0) as u64);
serde::Serialize::serialize(&stats, &serde_wasm_bindgen::Serializer::json_compatible())
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = closeIncoming)]
pub fn close_incoming(&self) {
self.inner.borrow_mut().close_incoming();
}
#[wasm_bindgen(js_name = closeOutgoing)]
pub fn close_outgoing(&self) {
self.inner.borrow_mut().close_outgoing();
}
pub fn close(&self) {
self.inner.borrow_mut().close();
}
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
fn emit_wasm_carrier_action(
handler: &Rc<RefCell<Option<js_sys::Function>>>,
action: serde_json::Value,
) {
let Some(handler) = handler.borrow().as_ref().cloned() else {
return;
};
let Ok(value) = serde::Serialize::serialize(
&action,
&serde_wasm_bindgen::Serializer::json_compatible(),
) else {
return;
};
if let Err(error) = handler.call1(&JsValue::UNDEFINED, &value) {
web_sys::console::error_2(
&JsValue::from_str("[OpenRTC][WASM carrier] action handler failed"),
&error,
);
}
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
fn wasm_carrier_failure_code(error: &anyhow::Error) -> &'static str {
let message = format!("{error:#}");
if message.contains("base generation changed before candidate acknowledgement")
|| message.contains("incumbent generation changed")
|| message.contains("replacement incumbent is stale")
{
"carrier-base-generation-stale"
} else if message.contains("authorization epoch changed") {
"carrier-authorization-stale"
} else if message.contains("stale before atomic commit")
|| message.contains("became stale during atomic commit")
{
"carrier-logical-generation-stale"
} else if message.contains("attempt was retired") || message.contains("retired upgrade") {
"carrier-attempt-retired"
} else {
"carrier-proof-failed"
}
}
struct TicketMeshStartGuard<'a> {
starts: &'a RefCell<HashSet<String>>,
id: String,
}
impl Drop for TicketMeshStartGuard<'_> {
fn drop(&mut self) {
self.starts.borrow_mut().remove(&self.id);
}
}
#[wasm_bindgen]
pub struct WasmClient {
inner: Arc<Client>,
portable_media: RefCell<crate::media::PortableMediaSession>,
broadcast_sessions: RefCell<HashMap<String, crate::broadcast::BroadcastSession>>,
broadcast_signers: RefCell<HashMap<String, crate::broadcast::BroadcastPublisherSigner>>,
ticket_meshes:
RefCell<HashMap<String, Rc<crate::ticket_mesh::wasm_session::WasmTicketMesh>>>,
ticket_mesh_starts: RefCell<HashSet<String>>,
identity_credential: Arc<Mutex<Option<String>>>,
last_auth_log: Arc<Mutex<Option<(bool, usize)>>>,
#[cfg(feature = "iroh-protocols-wasm")]
persistent_protocols:
Rc<tokio::sync::Mutex<Option<crate::wasm_docs_persistence::WasmPersistentDocsActor>>>,
#[cfg(feature = "managed-group-encryption")]
managed_group_controllers:
Rc<RefCell<HashMap<String, crate::managed_group_controller::ManagedGroupController>>>,
#[cfg(feature = "transport-webrtc")]
wasm_webrtc_carrier_sessions:
Rc<RefCell<HashMap<String, crate::wasm_webrtc_carrier::WasmWebRtcCarrierSession>>>,
#[cfg(feature = "transport-webrtc")]
wasm_webrtc_carrier_attempts: Rc<RefCell<HashMap<String, WasmWebRtcCarrierAttempt>>>,
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_carrier_action_handler: Rc<RefCell<Option<js_sys::Function>>>,
#[cfg(feature = "transport-moq")]
wasm_moq_carrier_sessions:
Rc<RefCell<HashMap<String, crate::wasm_moq_carrier::WasmMoqCarrierSession>>>,
#[cfg(feature = "transport-moq")]
wasm_moq_carrier_attempts: Rc<RefCell<HashMap<String, WasmMoqCarrierAttempt>>>,
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_remote_carrier_capabilities:
Rc<RefCell<HashMap<String, WasmRemoteCarrierCapabilities>>>,
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_carrier_peer_lifecycle: Rc<tokio::sync::Mutex<()>>,
}
#[cfg(feature = "transport-webrtc")]
async fn fail_wasm_webrtc_attempt(
inner: Arc<Client>,
attempts: Rc<RefCell<HashMap<String, WasmWebRtcCarrierAttempt>>>,
sessions: Rc<
RefCell<HashMap<String, crate::wasm_webrtc_carrier::WasmWebRtcCarrierSession>>,
>,
handler: Rc<RefCell<Option<js_sys::Function>>>,
attempt: WasmWebRtcCarrierAttempt,
failure_code: &'static str,
notify_peer: bool,
) {
let kind = crate::client::IrohPathKind::WebRtc;
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
if attempts
.borrow()
.get(&attempt.connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id)
{
attempts.borrow_mut().remove(&attempt.connection_id);
}
sessions.borrow_mut().remove(&attempt.bootstrap.upgrade_id);
if notify_peer {
if let Ok(failed) = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::failed_from(
&attempt.bootstrap,
failure_code,
) {
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "send-control",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": failed,
}),
);
}
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-webrtc",
"connectionId": attempt.connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": failure_code,
}),
);
let retry_pending = attempt.retry_count == 0
&& matches!(
failure_code,
"data-channel-failed"
| "ice-failed"
| "signaling-failed"
| "carrier-base-generation-stale"
);
if attempt.role == "initiator" && !retry_pending {
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "advance-carrier",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"failedRoute": "webrtc",
}),
);
}
}
#[cfg(feature = "transport-webrtc")]
async fn retire_selected_wasm_webrtc_carrier(
inner: Arc<Client>,
attempts: Rc<RefCell<HashMap<String, WasmWebRtcCarrierAttempt>>>,
sessions: Rc<
RefCell<HashMap<String, crate::wasm_webrtc_carrier::WasmWebRtcCarrierSession>>,
>,
handler: Rc<RefCell<Option<js_sys::Function>>>,
attempt: &WasmWebRtcCarrierAttempt,
failure_code: &'static str,
lifecycle_reason: &'static str,
) -> bool {
if !inner
.close_current_iroh_carrier_generation_with_reason(
&attempt.connection_id,
attempt.remote_endpoint_id,
crate::client::IrohPathKind::WebRtc,
attempt.generation.transport_generation.saturating_add(1),
lifecycle_reason,
)
.await
{
return false;
}
web_sys::console::warn_1(&JsValue::from_str(&format!(
"[OpenRTC][WebRTC carrier] selected mechanism ended connection_id={} upgrade_id={} failure_code={failure_code}",
attempt.connection_id, attempt.bootstrap.upgrade_id,
)));
if attempts
.borrow()
.get(&attempt.connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id)
{
attempts.borrow_mut().remove(&attempt.connection_id);
}
sessions.borrow_mut().remove(&attempt.bootstrap.upgrade_id);
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-webrtc",
"connectionId": attempt.connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": failure_code,
}),
);
true
}
#[cfg(feature = "transport-webrtc")]
impl WasmClient {
fn schedule_wasm_webrtc_carrier_watchdog(
&self,
bootstrap: crate::iroh_carrier_bootstrap::CarrierBootstrapFrame,
) {
let inner = self.inner.clone();
let attempts = self.wasm_webrtc_carrier_attempts.clone();
let sessions = self.wasm_webrtc_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
let mut completion_phase_observed = false;
loop {
let timeout = if completion_phase_observed { 45 } else { 30 };
gloo_timers::future::sleep(std::time::Duration::from_secs(timeout)).await;
let _lifecycle = lifecycle.lock().await;
let attempt = attempts
.borrow()
.values()
.find(|attempt| attempt.bootstrap.upgrade_id == bootstrap.upgrade_id)
.cloned();
let Some(attempt) = attempt else {
return;
};
if !inner
.wasm_carrier_upgrade_is_current(
&attempt.connection_id,
crate::client::IrohPathKind::WebRtc,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
if attempt.completion_started && !completion_phase_observed {
completion_phase_observed = true;
continue;
}
fail_wasm_webrtc_attempt(
inner,
attempts,
sessions,
handler,
attempt,
"carrier-timeout",
true,
)
.await;
return;
}
});
}
async fn fail_wasm_webrtc_carrier_attempt(
&self,
attempt: WasmWebRtcCarrierAttempt,
failure_code: &'static str,
notify_peer: bool,
) {
fail_wasm_webrtc_attempt(
self.inner.clone(),
self.wasm_webrtc_carrier_attempts.clone(),
self.wasm_webrtc_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
attempt,
failure_code,
notify_peer,
)
.await;
}
fn spawn_outbound_wasm_webrtc_carrier_completion(&self, attempt: WasmWebRtcCarrierAttempt) {
let inner = self.inner.clone();
let attempts = self.wasm_webrtc_carrier_attempts.clone();
let sessions = self.wasm_webrtc_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
let kind = crate::client::IrohPathKind::WebRtc;
let result = inner
.complete_outbound_wasm_carrier_upgrade(
&attempt.connection_id,
&attempt.remote_endpoint_id.to_string(),
&attempt.bootstrap.upgrade_id,
attempt.generation,
kind,
)
.await;
let _lifecycle = lifecycle.lock().await;
if let Err(error) = result {
let failure_code = wasm_carrier_failure_code(&error);
web_sys::console::error_1(&JsValue::from_str(&format!(
"[OpenRTC][WebRTC carrier] outbound candidate completion failed connection_id={} upgrade_id={} failure_code={} error={error:#}",
attempt.connection_id, attempt.bootstrap.upgrade_id, failure_code,
)));
fail_wasm_webrtc_attempt(
inner,
attempts,
sessions,
handler,
attempt,
failure_code,
true,
)
.await;
return;
}
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "selected",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"upgradeId": attempt.bootstrap.upgrade_id,
"family": "iroh",
"carrier": "webrtc",
"transportGeneration": attempt.generation.transport_generation.saturating_add(1),
"routeGeneration": 0,
}),
);
});
}
fn take_ready_outbound_wasm_webrtc_carrier_attempt(
&self,
connection_id: &str,
upgrade_id: &str,
) -> Option<WasmWebRtcCarrierAttempt> {
if !self
.wasm_webrtc_carrier_sessions
.borrow()
.contains_key(upgrade_id)
{
return None;
}
let mut attempts = self.wasm_webrtc_carrier_attempts.borrow_mut();
let attempt = attempts.get_mut(connection_id).filter(|attempt| {
attempt.bootstrap.upgrade_id == upgrade_id
&& attempt.role == "initiator"
&& attempt.remote_ready
&& !attempt.completion_started
})?;
attempt.completion_started = true;
Some(attempt.clone())
}
fn spawn_inbound_wasm_webrtc_carrier_completion(&self, attempt: WasmWebRtcCarrierAttempt) {
let inner = self.inner.clone();
let attempts = self.wasm_webrtc_carrier_attempts.clone();
let sessions = self.wasm_webrtc_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
let kind = crate::client::IrohPathKind::WebRtc;
let node = inner.iroh_node.read().await.as_ref().cloned();
let result = async {
let node = node.ok_or_else(|| anyhow::anyhow!("Iroh node is unavailable"))?;
let authorization_expires_at_ms = attempt
.inbound_authorization_expires_at_ms
.ok_or_else(|| anyhow::anyhow!("browser WebRTC inbound authorization is missing"))?;
let candidate = node
.wait_for_inbound_replacement_candidate(
attempt.remote_endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID,
authorization_expires_at_ms,
std::time::Duration::from_secs(40),
)
.await?;
let policy_epoch = inner
.require_current_wasm_carrier_upgrade_fence(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await?;
let authorization_fence = inner
.capture_wasm_iroh_carrier_authorization_fence(
&attempt.connection_id,
attempt.generation,
kind,
policy_epoch,
)?;
let proof = inner.wasm_candidate_proof_probe(
&attempt.connection_id,
&attempt.bootstrap.upgrade_id,
attempt.bootstrap.base.transport_generation,
attempt.bootstrap.base.route_generation,
kind,
crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID,
)?;
let (send, probe, ack) = inner
.receive_inbound_wasm_carrier_candidate_proof(
&attempt.connection_id,
&candidate,
&proof,
)
.await?;
anyhow::ensure!(
inner
.current_wasm_peer_data_generation(&attempt.connection_id, None)
.await
== Some(attempt.generation),
"browser WebRTC inbound carrier base generation changed before candidate acknowledgement"
);
inner
.send_inbound_wasm_carrier_candidate_ack(send, &ack)
.await?;
let (commit_send, committed) = inner
.receive_inbound_wasm_carrier_commit(
&attempt.connection_id,
&candidate,
&probe,
)
.await?;
let committed_replacement = inner
.commit_proven_wasm_carrier_candidate(
&attempt.connection_id,
&attempt.bootstrap.upgrade_id,
attempt.generation,
kind,
authorization_fence,
candidate,
)
.await?;
inner
.send_inbound_wasm_carrier_candidate_ack(commit_send, &committed)
.await?;
let replacement_transport_stable_id = committed_replacement
.logical_result()
.transport_stable_id
.ok_or_else(|| {
anyhow::anyhow!(
"browser WebRTC replacement has no stable ID"
)
})?;
committed_replacement.finish(b"wasm-custom-transport-upgrade");
inner
.publish_committed_wasm_carrier_route(
&attempt.connection_id,
kind,
replacement_transport_stable_id,
)
.await?;
Ok::<(), anyhow::Error>(())
}
.await;
if let Some(expires_at_ms) = attempt.inbound_authorization_expires_at_ms {
if let Some(node) = inner.iroh_node.read().await.as_ref().cloned() {
node.revoke_inbound_replacement_if_current(
attempt.remote_endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID,
expires_at_ms,
)
.await;
}
}
let _lifecycle = lifecycle.lock().await;
if let Err(error) = result {
let failure_code = wasm_carrier_failure_code(&error);
web_sys::console::error_1(&JsValue::from_str(&format!(
"[OpenRTC][WebRTC carrier] inbound candidate completion failed connection_id={} upgrade_id={} failure_code={} error={error:#}",
attempt.connection_id, attempt.bootstrap.upgrade_id, failure_code,
)));
fail_wasm_webrtc_attempt(
inner,
attempts,
sessions,
handler,
attempt,
failure_code,
true,
)
.await;
return;
}
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "selected",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"upgradeId": attempt.bootstrap.upgrade_id,
"family": "iroh",
"carrier": "webrtc",
"transportGeneration": attempt.generation.transport_generation.saturating_add(1),
"routeGeneration": 0,
}),
);
});
}
}
#[cfg(feature = "transport-moq")]
fn wasm_moq_carrier_namespaces(
local_endpoint_id: &str,
remote_endpoint_id: &str,
carrier_session_id: &str,
) -> (String, String, &'static str) {
let (first, second) = if local_endpoint_id <= remote_endpoint_id {
(local_endpoint_id, remote_endpoint_id)
} else {
(remote_endpoint_id, local_endpoint_id)
};
let base = format!("openrtc/iroh-carrier/moq/{first}/{second}/{carrier_session_id}");
(
format!("{base}/from/{local_endpoint_id}"),
format!("{base}/from/{remote_endpoint_id}"),
"iroh-packets",
)
}
#[cfg(feature = "transport-moq")]
async fn fail_wasm_moq_attempt(
inner: Arc<Client>,
attempts: Rc<RefCell<HashMap<String, WasmMoqCarrierAttempt>>>,
sessions: Rc<RefCell<HashMap<String, crate::wasm_moq_carrier::WasmMoqCarrierSession>>>,
handler: Rc<RefCell<Option<js_sys::Function>>>,
attempt: WasmMoqCarrierAttempt,
failure_code: &'static str,
notify_peer: bool,
) {
let kind = crate::client::IrohPathKind::Moq;
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
if attempts
.borrow()
.get(&attempt.connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id)
{
attempts.borrow_mut().remove(&attempt.connection_id);
}
sessions.borrow_mut().remove(&attempt.bootstrap.upgrade_id);
if notify_peer {
if let Ok(failed) = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::failed_from(
&attempt.bootstrap,
failure_code,
) {
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "send-control",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": failed,
}),
);
}
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-moq",
"connectionId": attempt.connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": failure_code,
}),
);
let retry_pending = attempt.role == "initiator"
&& attempt.retry_count == 0
&& failure_code == "carrier-base-generation-stale";
if attempt.role == "initiator" && !retry_pending {
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "advance-carrier",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"failedRoute": "moq",
}),
);
}
}
#[cfg(feature = "transport-moq")]
async fn retire_selected_wasm_moq_carrier(
inner: Arc<Client>,
attempts: Rc<RefCell<HashMap<String, WasmMoqCarrierAttempt>>>,
sessions: Rc<RefCell<HashMap<String, crate::wasm_moq_carrier::WasmMoqCarrierSession>>>,
handler: Rc<RefCell<Option<js_sys::Function>>>,
attempt: &WasmMoqCarrierAttempt,
failure_code: &'static str,
lifecycle_reason: &'static str,
) -> bool {
if !inner
.close_current_iroh_carrier_generation_with_reason(
&attempt.connection_id,
attempt.remote_endpoint_id,
crate::client::IrohPathKind::Moq,
attempt.generation.transport_generation.saturating_add(1),
lifecycle_reason,
)
.await
{
return false;
}
web_sys::console::warn_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier] selected mechanism ended connection_id={} upgrade_id={} failure_code={failure_code}",
attempt.connection_id, attempt.bootstrap.upgrade_id,
)));
if attempts
.borrow()
.get(&attempt.connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id)
{
attempts.borrow_mut().remove(&attempt.connection_id);
}
sessions.borrow_mut().remove(&attempt.bootstrap.upgrade_id);
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-moq",
"connectionId": attempt.connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": failure_code,
}),
);
true
}
#[cfg(feature = "transport-moq")]
impl WasmClient {
fn schedule_wasm_moq_carrier_watchdog(
&self,
bootstrap: crate::iroh_carrier_bootstrap::CarrierBootstrapFrame,
) {
let inner = self.inner.clone();
let attempts = self.wasm_moq_carrier_attempts.clone();
let sessions = self.wasm_moq_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
gloo_timers::future::sleep(std::time::Duration::from_secs(30)).await;
let _lifecycle = lifecycle.lock().await;
let attempt = attempts
.borrow()
.values()
.find(|attempt| attempt.bootstrap.upgrade_id == bootstrap.upgrade_id)
.cloned();
if let Some(attempt) = attempt {
if !inner
.wasm_carrier_upgrade_is_current(
&attempt.connection_id,
crate::client::IrohPathKind::Moq,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
fail_wasm_moq_attempt(
inner,
attempts,
sessions,
handler,
attempt,
"carrier-timeout",
true,
)
.await;
}
});
}
async fn fail_wasm_moq_carrier_attempt(
&self,
attempt: WasmMoqCarrierAttempt,
failure_code: &'static str,
notify_peer: bool,
) {
fail_wasm_moq_attempt(
self.inner.clone(),
self.wasm_moq_carrier_attempts.clone(),
self.wasm_moq_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
attempt,
failure_code,
notify_peer,
)
.await;
}
fn spawn_outbound_wasm_moq_carrier_completion(&self, attempt: WasmMoqCarrierAttempt) {
let inner = self.inner.clone();
let attempts = self.wasm_moq_carrier_attempts.clone();
let sessions = self.wasm_moq_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
let kind = crate::client::IrohPathKind::Moq;
let result = inner
.complete_outbound_wasm_carrier_upgrade(
&attempt.connection_id,
&attempt.remote_endpoint_id.to_string(),
&attempt.bootstrap.upgrade_id,
attempt.generation,
kind,
)
.await;
let _lifecycle = lifecycle.lock().await;
if let Err(error) = result {
let failure_code = wasm_carrier_failure_code(&error);
web_sys::console::error_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier] outbound candidate completion failed connection_id={} upgrade_id={} failure_code={} error={error:#}",
attempt.connection_id, attempt.bootstrap.upgrade_id, failure_code,
)));
fail_wasm_moq_attempt(
inner,
attempts,
sessions,
handler,
attempt,
failure_code,
true,
)
.await;
return;
}
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "selected",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"upgradeId": attempt.bootstrap.upgrade_id,
"family": "iroh",
"carrier": "moq",
"transportGeneration": attempt.generation.transport_generation.saturating_add(1),
"routeGeneration": 0,
}),
);
});
}
fn spawn_inbound_wasm_moq_carrier_completion(&self, attempt: WasmMoqCarrierAttempt) {
let inner = self.inner.clone();
let attempts = self.wasm_moq_carrier_attempts.clone();
let sessions = self.wasm_moq_carrier_sessions.clone();
let handler = self.wasm_carrier_action_handler.clone();
let lifecycle = self.wasm_carrier_peer_lifecycle.clone();
spawn_local(async move {
let kind = crate::client::IrohPathKind::Moq;
let node = inner.iroh_node.read().await.as_ref().cloned();
let result = async {
let node = node.ok_or_else(|| anyhow::anyhow!("Iroh node is unavailable"))?;
let authorization_expires_at_ms = attempt
.inbound_authorization_expires_at_ms
.ok_or_else(|| anyhow::anyhow!("browser MoQ inbound authorization is missing"))?;
let candidate = node
.wait_for_inbound_replacement_candidate(
attempt.remote_endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID,
authorization_expires_at_ms,
std::time::Duration::from_secs(40),
)
.await?;
let policy_epoch = inner
.require_current_wasm_carrier_upgrade_fence(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await?;
let authorization_fence = inner
.capture_wasm_iroh_carrier_authorization_fence(
&attempt.connection_id,
attempt.generation,
kind,
policy_epoch,
)?;
let proof = inner.wasm_candidate_proof_probe(
&attempt.connection_id,
&attempt.bootstrap.upgrade_id,
attempt.bootstrap.base.transport_generation,
attempt.bootstrap.base.route_generation,
kind,
crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID,
)?;
let (send, probe, ack) = inner
.receive_inbound_wasm_carrier_candidate_proof(
&attempt.connection_id,
&candidate,
&proof,
)
.await?;
anyhow::ensure!(
inner
.current_wasm_peer_data_generation(&attempt.connection_id, None)
.await
== Some(attempt.generation),
"browser MoQ inbound carrier base generation changed before candidate acknowledgement"
);
inner
.send_inbound_wasm_carrier_candidate_ack(send, &ack)
.await?;
let (commit_send, committed) = inner
.receive_inbound_wasm_carrier_commit(
&attempt.connection_id,
&candidate,
&probe,
)
.await?;
let committed_replacement = inner
.commit_proven_wasm_carrier_candidate(
&attempt.connection_id,
&attempt.bootstrap.upgrade_id,
attempt.generation,
kind,
authorization_fence,
candidate,
)
.await?;
inner
.send_inbound_wasm_carrier_candidate_ack(commit_send, &committed)
.await?;
let replacement_transport_stable_id = committed_replacement
.logical_result()
.transport_stable_id
.ok_or_else(|| {
anyhow::anyhow!("browser MoQ replacement has no stable ID")
})?;
committed_replacement.finish(b"wasm-custom-transport-upgrade");
inner
.publish_committed_wasm_carrier_route(
&attempt.connection_id,
kind,
replacement_transport_stable_id,
)
.await?;
Ok::<(), anyhow::Error>(())
}
.await;
if let Some(expires_at_ms) = attempt.inbound_authorization_expires_at_ms {
if let Some(node) = inner.iroh_node.read().await.as_ref().cloned() {
node.revoke_inbound_replacement_if_current(
attempt.remote_endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID,
expires_at_ms,
)
.await;
}
}
let _lifecycle = lifecycle.lock().await;
if let Err(error) = result {
let failure_code = wasm_carrier_failure_code(&error);
web_sys::console::error_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier] inbound candidate completion failed connection_id={} upgrade_id={} failure_code={} error={error:#}",
attempt.connection_id, attempt.bootstrap.upgrade_id, failure_code,
)));
fail_wasm_moq_attempt(
inner,
attempts,
sessions,
handler,
attempt,
failure_code,
true,
)
.await;
return;
}
if !inner
.retire_wasm_carrier_upgrade(
&attempt.connection_id,
kind,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await
{
return;
}
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "selected",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"upgradeId": attempt.bootstrap.upgrade_id,
"family": "iroh",
"carrier": "moq",
"transportGeneration": attempt.generation.transport_generation.saturating_add(1),
"routeGeneration": 0,
}),
);
});
}
async fn handle_wasm_moq_bootstrap(
&self,
connection_id: String,
remote_endpoint_id: String,
bootstrap: crate::iroh_carrier_bootstrap::CarrierBootstrapFrame,
) -> Result<bool, JsValue> {
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let kind = crate::client::IrohPathKind::Moq;
match bootstrap.action {
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Request => {
let local_endpoint_id = self
.inner
.current_node_id()
.await
.ok_or_else(|| JsValue::from_str("local endpoint id is unavailable"))?;
if local_endpoint_id.as_str() >= remote_endpoint_id.as_str()
|| !self.inner.is_moq_carrier_enabled().await
{
return Ok(true);
}
if !crate::iroh_connection_policy::custom_carrier_base_allows(
self.inner.iroh_path_kind(&remote_endpoint_id).await,
crate::client::IrohPathKind::Moq,
false,
) {
return Ok(true);
}
let base_connection = self
.inner
.get_connection(endpoint_id)
.await
.ok_or_else(|| JsValue::from_str("MoQ carrier base is unavailable"))?;
let generation = self
.inner
.current_wasm_peer_data_generation(&connection_id, None)
.await
.ok_or_else(|| {
JsValue::from_str("MoQ carrier generation is unavailable")
})?;
if crate::transport_generation::for_connection(&base_connection)
!= generation.transport_stable_id
{
let failed =
crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::failed_from(
&bootstrap,
"carrier-base-generation-stale",
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"envelope": failed,
}),
);
return Ok(true);
}
if !matches!(
self.inner
.reserve_wasm_carrier_upgrade(
&connection_id,
kind,
&bootstrap.upgrade_id,
generation,
)
.await,
crate::client::WasmCarrierUpgradeReservation::Reserved
) {
return Ok(true);
}
let node = self
.inner
.iroh_node
.read()
.await
.as_ref()
.cloned()
.ok_or_else(|| JsValue::from_str("Iroh node is unavailable"))?;
let authorization_expiry = node
.authorize_pending_inbound_replacement(
endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID,
std::time::Duration::from_secs(45),
)
.await;
let previous = self.wasm_moq_carrier_attempts.borrow_mut().insert(
connection_id.clone(),
WasmMoqCarrierAttempt {
connection_id: connection_id.clone(),
remote_endpoint_id: endpoint_id,
bootstrap: bootstrap.clone(),
generation,
role: "responder",
prepared: false,
retry_sent: false,
retry_count: bootstrap.attempt.saturating_sub(1),
inbound_authorization_expires_at_ms: Some(authorization_expiry),
},
);
if let Some(previous) = previous {
self.wasm_moq_carrier_sessions
.borrow_mut()
.remove(&previous.bootstrap.upgrade_id);
}
let (publish_namespace, subscribe_namespace, track_name) =
wasm_moq_carrier_namespaces(
&local_endpoint_id,
&remote_endpoint_id,
&bootstrap.carrier_session_id,
);
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "prepare-moq",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"role": "responder",
"upgradeId": bootstrap.upgrade_id,
"carrierSessionId": bootstrap.carrier_session_id,
"transportGeneration": generation.transport_generation.saturating_add(1),
"publishNamespace": publish_namespace,
"subscribeNamespace": subscribe_namespace,
"trackName": track_name,
}),
);
self.schedule_wasm_moq_carrier_watchdog(bootstrap);
}
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Ready => {
let attempt = self
.wasm_moq_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| bootstrap.is_response_to(&attempt.bootstrap))
.cloned();
if let Some(attempt) = attempt {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "activate-moq",
"connectionId": attempt.connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
}),
);
}
}
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Failed => {
let attempt = self
.wasm_moq_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| bootstrap.is_response_to(&attempt.bootstrap))
.cloned();
if let Some(attempt) = attempt {
let failure_code = crate::client::wasm_peer_carrier_failure_code(
bootstrap.failure_code.as_deref(),
);
if crate::client::should_rearm_wasm_carrier_event_retry(
failure_code,
attempt.role == "initiator",
false,
) {
if let Some(current) = self
.wasm_moq_carrier_attempts
.borrow_mut()
.get_mut(&connection_id)
.filter(|current| {
current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id
})
{
current.retry_sent = false;
}
return Ok(true);
}
let should_retry = attempt.role == "initiator"
&& attempt.retry_count == 0
&& failure_code == "carrier-base-generation-stale";
let remote_endpoint_id = attempt.remote_endpoint_id.to_string();
self.fail_wasm_moq_carrier_attempt(attempt, failure_code, false)
.await;
if should_retry {
gloo_timers::future::sleep(std::time::Duration::from_millis(500)).await;
let _ = self
.begin_iroh_moq_carrier_attempt(
connection_id,
remote_endpoint_id,
1,
)
.await;
}
}
}
}
Ok(true)
}
}
#[wasm_bindgen]
impl WasmClient {
#[wasm_bindgen(constructor)]
pub fn new(api_key: String) -> Result<WasmClient, JsValue> {
let api_key = crate::validate_api_key(&api_key)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let app_tag = crate::app_tag_from_api_key(api_key);
let identity_credential: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
let credential_state = identity_credential.clone();
let identity_credential_provider = Box::new(move || -> Option<String> {
credential_state.lock().ok().and_then(|guard| guard.clone())
});
Ok(Self {
inner: Arc::new(Client::new_provider_neutral(
app_tag,
identity_credential_provider,
)),
portable_media: RefCell::new(crate::media::PortableMediaSession::default()),
broadcast_sessions: RefCell::new(HashMap::new()),
broadcast_signers: RefCell::new(HashMap::new()),
ticket_meshes: RefCell::new(HashMap::new()),
ticket_mesh_starts: RefCell::new(HashSet::new()),
identity_credential,
last_auth_log: Arc::new(Mutex::new(None)),
#[cfg(feature = "iroh-protocols-wasm")]
persistent_protocols: Rc::new(tokio::sync::Mutex::new(None)),
#[cfg(feature = "managed-group-encryption")]
managed_group_controllers: Rc::new(RefCell::new(HashMap::new())),
#[cfg(feature = "transport-webrtc")]
wasm_webrtc_carrier_sessions: Rc::new(RefCell::new(HashMap::new())),
#[cfg(feature = "transport-webrtc")]
wasm_webrtc_carrier_attempts: Rc::new(RefCell::new(HashMap::new())),
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_carrier_action_handler: Rc::new(RefCell::new(None)),
#[cfg(feature = "transport-moq")]
wasm_moq_carrier_sessions: Rc::new(RefCell::new(HashMap::new())),
#[cfg(feature = "transport-moq")]
wasm_moq_carrier_attempts: Rc::new(RefCell::new(HashMap::new())),
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_remote_carrier_capabilities: Rc::new(RefCell::new(HashMap::new())),
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
wasm_carrier_peer_lifecycle: Rc::new(tokio::sync::Mutex::new(())),
})
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __initManagedRoomGroup)]
pub fn init_managed_room_group(
&self,
avenue_key: String,
device_id: String,
wrapping_key: Vec<u8>,
sealed_state: Option<Vec<u8>>,
) -> Result<JsValue, JsValue> {
let wrapping_key: [u8; 32] = wrapping_key
.try_into()
.map_err(|_| JsValue::from_str("managed room wrapping key must be 32 bytes"))?;
let mut controllers = self.managed_group_controllers.borrow_mut();
if !controllers.contains_key(&avenue_key) {
let controller = crate::managed_group_controller::ManagedGroupController::new(
device_id,
wrapping_key,
sealed_state.as_deref(),
)
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
controllers.insert(avenue_key.clone(), controller);
}
let controller = controllers
.get(&avenue_key)
.ok_or_else(|| JsValue::from_str("managed room group initialization failed"))?;
let action = controller
.publish_key_package()
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
serde::Serialize::serialize(&action, &serde_wasm_bindgen::Serializer::json_compatible())
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __handleManagedRoomPreparePage)]
pub fn handle_managed_room_prepare_page(
&self,
avenue_key: String,
page: JsValue,
) -> Result<JsValue, JsValue> {
let page = serde_wasm_bindgen::from_value(page)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let mut controllers = self.managed_group_controllers.borrow_mut();
let controller = controllers
.get_mut(&avenue_key)
.ok_or_else(|| JsValue::from_str("managed room group is not initialized"))?;
let actions = controller
.handle_prepare_page(page)
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
serde::Serialize::serialize(
&actions,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __handleManagedRoomArtifactChunk)]
pub fn handle_managed_room_artifact_chunk(
&self,
avenue_key: String,
chunk: JsValue,
) -> Result<JsValue, JsValue> {
let chunk = serde_wasm_bindgen::from_value(chunk)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let mut controllers = self.managed_group_controllers.borrow_mut();
let controller = controllers
.get_mut(&avenue_key)
.ok_or_else(|| JsValue::from_str("managed room group is not initialized"))?;
let actions = controller
.handle_artifact_chunk(chunk)
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
serde::Serialize::serialize(
&actions,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __sealManagedRoomPayload)]
#[allow(clippy::too_many_arguments)]
pub fn seal_managed_room_payload(
&self,
avenue_key: String,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: String,
channel: String,
priority: u8,
zone_id: Option<String>,
payload: Vec<u8>,
) -> Result<JsValue, JsValue> {
let protected = self
.managed_group_controllers
.borrow_mut()
.get_mut(&avenue_key)
.ok_or_else(|| JsValue::from_str("managed room group is not initialized"))?
.seal_payload(
architecture_epoch,
encryption_epoch,
&message_id,
&channel,
priority,
zone_id.as_deref(),
&payload,
)
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
serde_wasm_bindgen::to_value(&serde_json::json!({
"payload": base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
protected.data,
),
"sealedState": base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
protected.sealed_state,
),
}))
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __openManagedRoomPayload)]
#[allow(clippy::too_many_arguments)]
pub fn open_managed_room_payload(
&self,
avenue_key: String,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: String,
channel: String,
priority: u8,
zone_id: Option<String>,
ciphertext: Vec<u8>,
) -> Result<JsValue, JsValue> {
let protected = self
.managed_group_controllers
.borrow_mut()
.get_mut(&avenue_key)
.ok_or_else(|| JsValue::from_str("managed room group is not initialized"))?
.open_payload(
architecture_epoch,
encryption_epoch,
&message_id,
&channel,
priority,
zone_id.as_deref(),
&ciphertext,
)
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
serde_wasm_bindgen::to_value(&serde_json::json!({
"payload": base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
protected.data,
),
"sealedState": base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
protected.sealed_state,
),
}))
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "managed-group-encryption")]
#[wasm_bindgen(js_name = __forgetManagedRoomGroup)]
pub fn forget_managed_room_group(&self, avenue_key: String) {
self.managed_group_controllers
.borrow_mut()
.remove(&avenue_key);
}
#[wasm_bindgen(js_name = __prepareOpenRtcBroadcastGrant)]
pub fn prepare_openrtc_broadcast_grant(
&self,
grant_token: String,
issuer_public_key: Vec<u8>,
now_ms: f64,
) -> Result<JsValue, JsValue> {
let issuer_public_key: [u8; 32] = issuer_public_key
.try_into()
.map_err(|_| JsValue::from_str("broadcast issuer key must be 32 bytes"))?;
let issuer = ed25519_dalek::VerifyingKey::from_bytes(&issuer_public_key)
.map_err(|_| JsValue::from_str("broadcast issuer key is invalid"))?;
if !now_ms.is_finite()
|| now_ms < 0.0
|| now_ms.fract() != 0.0
|| now_ms > crate::media::MAX_JAVASCRIPT_SAFE_INTEGER as f64
{
return Err(JsValue::from_str(
"broadcast now must be a non-negative JavaScript-safe integer",
));
}
let challenge = self
.inner
.broadcasts()
.prepare_grant_verification(&grant_token, &issuer, now_ms as u64)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let object = js_sys::Object::new();
js_sys::Reflect::set(
&object,
&JsValue::from_str("handle"),
&JsValue::from_str(&challenge.handle),
)?;
js_sys::Reflect::set(
&object,
&JsValue::from_str("signingBytes"),
&js_sys::Uint8Array::from(challenge.signing_bytes.as_slice()),
)?;
js_sys::Reflect::set(
&object,
&JsValue::from_str("expiresAtMs"),
&JsValue::from_f64(challenge.expires_at_ms as f64),
)?;
Ok(object.into())
}
#[wasm_bindgen(js_name = __openOpenRtcBroadcast)]
pub fn open_openrtc_broadcast(
&self,
grant_token: String,
issuer_public_key: Vec<u8>,
binding_challenge_handle: String,
binding_signature: Vec<u8>,
publication_signer_handle: Option<String>,
now_ms: f64,
) -> Result<JsValue, JsValue> {
let issuer_public_key: [u8; 32] = issuer_public_key
.try_into()
.map_err(|_| JsValue::from_str("broadcast issuer key must be 32 bytes"))?;
let issuer = ed25519_dalek::VerifyingKey::from_bytes(&issuer_public_key)
.map_err(|_| JsValue::from_str("broadcast issuer key is invalid"))?;
if !now_ms.is_finite()
|| now_ms < 0.0
|| now_ms.fract() != 0.0
|| now_ms > crate::media::MAX_JAVASCRIPT_SAFE_INTEGER as f64
{
return Err(JsValue::from_str(
"broadcast now must be a non-negative JavaScript-safe integer",
));
}
let now_ms = now_ms as u64;
let grant = self
.inner
.broadcasts()
.complete_grant_verification(
&grant_token,
&issuer,
&binding_challenge_handle,
&binding_signature,
now_ms,
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let grant_generation = grant.grant_generation();
let session = match publication_signer_handle {
Some(handle) => {
let signers = self.broadcast_signers.borrow();
let signer = signers
.get(&handle)
.ok_or_else(|| JsValue::from_str("broadcast signer handle is invalid"))?;
self.inner
.broadcasts()
.open_publisher(grant, signer, now_ms)
}
None => self.inner.broadcasts().open(grant, now_ms),
}
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let handle_id = format!("{}:{grant_generation}", session.id());
let result = serde_json::json!({
"id": session.id(),
"handleId": handle_id,
"role": session.role(),
"state": session.state(),
"budget": session.budget(),
"sourceSlots": session.source_slots(),
});
self.broadcast_sessions
.borrow_mut()
.insert(handle_id, session);
serde::Serialize::serialize(&result, &serde_wasm_bindgen::Serializer::json_compatible())
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __createOpenRtcBroadcastPublisherSigner)]
pub fn create_openrtc_broadcast_publisher_signer(&self) -> Result<JsValue, JsValue> {
let signer = crate::broadcast::BroadcastPublisherSigner::generate()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let mut handle_bytes = [0_u8; 16];
getrandom::getrandom(&mut handle_bytes)
.map_err(|_| JsValue::from_str("broadcast signer handle entropy failed"))?;
let handle = hex::encode(handle_bytes);
let public_key = signer.verifying_key();
let object = js_sys::Object::new();
js_sys::Reflect::set(
&object,
&JsValue::from_str("handle"),
&JsValue::from_str(&handle),
)?;
js_sys::Reflect::set(
&object,
&JsValue::from_str("publicKey"),
&js_sys::Uint8Array::from(public_key.as_slice()),
)?;
self.broadcast_signers.borrow_mut().insert(handle, signer);
Ok(object.into())
}
#[wasm_bindgen(js_name = __releaseOpenRtcBroadcastPublisherSigner)]
pub fn release_openrtc_broadcast_publisher_signer(&self, handle: String) -> bool {
self.broadcast_signers
.borrow_mut()
.remove(&handle)
.is_some()
}
#[wasm_bindgen(js_name = __takeOpenRtcBroadcastActions)]
pub fn take_openrtc_broadcast_actions(
&self,
handle_id: String,
max: usize,
) -> Result<JsValue, JsValue> {
let sessions = self.broadcast_sessions.borrow();
let session = sessions
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?;
serde::Serialize::serialize(
&session.take_actions(max),
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __getOpenRtcBroadcastStats)]
pub fn get_openrtc_broadcast_stats(&self, handle_id: String) -> Result<JsValue, JsValue> {
let sessions = self.broadcast_sessions.borrow();
let stats = sessions
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.stats();
serde::Serialize::serialize(&stats, &serde_wasm_bindgen::Serializer::json_compatible())
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __observeOpenRtcBroadcastAdapter)]
pub fn observe_openrtc_broadcast_adapter(
&self,
handle_id: String,
observation: JsValue,
) -> Result<(), JsValue> {
let observation = serde_wasm_bindgen::from_value(observation)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.observe(observation)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __closeOpenRtcBroadcast)]
pub fn close_openrtc_broadcast(&self, handle_id: String) -> Result<(), JsValue> {
let session = self
.broadcast_sessions
.borrow_mut()
.remove(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?;
session.close();
Ok(())
}
#[wasm_bindgen(js_name = __beginOpenRtcBroadcastPublication)]
#[allow(clippy::too_many_arguments)]
pub fn begin_openrtc_broadcast_publication(
&self,
handle_id: String,
source_slot: String,
publication_id: Option<String>,
kind: String,
codec: String,
clock_rate: u32,
coded_width: Option<u32>,
coded_height: Option<u32>,
channels: Option<u16>,
) -> Result<JsValue, JsValue> {
let publication_id = match publication_id {
Some(value) if !value.is_empty() => {
value
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?
}
_ => {
let mut bytes = [0_u8; 16];
getrandom::getrandom(&mut bytes)
.map_err(|_| JsValue::from_str("media publication entropy failed"))?;
crate::media::PublicationId::from_bytes(bytes)
}
};
let kind: crate::media::MediaKind =
kind.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
if !matches!(
kind,
crate::media::MediaKind::Audio | crate::media::MediaKind::Video
) {
return Err(JsValue::from_str(
"browser broadcast publications must be audio or video",
));
}
let codec: crate::media::MediaCodec =
codec
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
let publication = self
.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.begin_publication(
&source_slot,
crate::media::MediaPublicationConfig {
publication_id,
media_generation: 0,
kind,
codec,
clock_rate,
coded_width,
coded_height,
channels,
},
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
serde::Serialize::serialize(
&serde_json::json!({
"publicationId": publication.publication_id.to_string(),
"mediaGeneration": publication.media_generation,
"control": Vec::<u8>::new(),
}),
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __publishOpenRtcBroadcastSample)]
#[allow(clippy::too_many_arguments)]
pub fn publish_openrtc_broadcast_sample(
&self,
handle_id: String,
publication_id: String,
timestamp_us: f64,
duration_us: u32,
keyframe: bool,
discardable: bool,
payload: Vec<u8>,
) -> Result<(), JsValue> {
if !timestamp_us.is_finite()
|| timestamp_us < 0.0
|| timestamp_us.fract() != 0.0
|| timestamp_us > crate::media::MAX_JAVASCRIPT_SAFE_INTEGER as f64
{
return Err(JsValue::from_str(
"media timestamp must be a non-negative JavaScript-safe integer",
));
}
let publication_id =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.publish(
publication_id,
crate::media::EncodedMediaSample {
timestamp_us: timestamp_us as u64,
duration_us,
keyframe,
discardable,
payload,
},
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __retireOpenRtcBroadcastPublication)]
pub fn retire_openrtc_broadcast_publication(
&self,
handle_id: String,
publication_id: String,
) -> Result<(), JsValue> {
let publication_id =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.retire_publication(publication_id)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __pauseOpenRtcBroadcastPublication)]
pub fn pause_openrtc_broadcast_publication(
&self,
handle_id: String,
publication_id: String,
) -> Result<(), JsValue> {
let publication_id =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.pause_publication(publication_id)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __acceptOpenRtcBroadcastObject)]
pub fn accept_openrtc_broadcast_object(
&self,
handle_id: String,
source_slot: String,
encoded: Vec<u8>,
) -> Result<JsValue, JsValue> {
let accepted = self
.broadcast_sessions
.borrow()
.get(&handle_id)
.ok_or_else(|| JsValue::from_str("broadcast session is unavailable"))?
.accept_media_for_source(&source_slot, &encoded)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
serde::Serialize::serialize(
&accepted,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __beginOpenRtcMediaPublication)]
#[allow(clippy::too_many_arguments)]
pub fn begin_openrtc_media_publication(
&self,
publication_id: Option<String>,
kind: String,
codec: String,
clock_rate: u32,
coded_width: Option<u32>,
coded_height: Option<u32>,
channels: Option<u16>,
) -> Result<JsValue, JsValue> {
let publication_id = match publication_id {
Some(value) if !value.is_empty() => {
value
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?
}
_ => {
let mut bytes = [0_u8; 16];
getrandom::getrandom(&mut bytes)
.map_err(|_| JsValue::from_str("media publication entropy failed"))?;
crate::media::PublicationId::from_bytes(bytes)
}
};
let kind: crate::media::MediaKind =
kind.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
if !matches!(
kind,
crate::media::MediaKind::Audio | crate::media::MediaKind::Video
) {
return Err(JsValue::from_str(
"browser media publications must be audio or video",
));
}
let codec: crate::media::MediaCodec =
codec
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
let publication = crate::media::MediaPublicationConfig {
publication_id,
media_generation: 0,
kind,
codec,
clock_rate,
coded_width,
coded_height,
channels,
};
let (publication, control) = self
.portable_media
.borrow_mut()
.begin_publication(publication)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let object = js_sys::Object::new();
js_sys::Reflect::set(
&object,
&JsValue::from_str("publicationId"),
&JsValue::from_str(&publication_id.to_string()),
)?;
js_sys::Reflect::set(
&object,
&JsValue::from_str("mediaGeneration"),
&JsValue::from_f64(publication.media_generation as f64),
)?;
js_sys::Reflect::set(
&object,
&JsValue::from_str("control"),
&js_sys::Uint8Array::from(control.as_slice()),
)?;
Ok(object.into())
}
#[wasm_bindgen(js_name = __encodeOpenRtcMediaSample)]
pub fn encode_openrtc_media_sample(
&self,
publication_id: String,
timestamp_us: f64,
duration_us: u32,
keyframe: bool,
discardable: bool,
payload: Vec<u8>,
) -> Result<Vec<u8>, JsValue> {
fn safe_u64(value: f64, field: &str) -> Result<u64, JsValue> {
if !value.is_finite()
|| value < 0.0
|| value.fract() != 0.0
|| value > crate::media::MAX_JAVASCRIPT_SAFE_INTEGER as f64
{
return Err(JsValue::from_str(&format!(
"media {field} must be a non-negative JavaScript-safe integer"
)));
}
Ok(value as u64)
}
let publication_id: crate::media::PublicationId =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.portable_media
.borrow_mut()
.encode_sample(
publication_id,
crate::media::EncodedMediaSample {
timestamp_us: safe_u64(timestamp_us, "timestamp")?,
duration_us,
keyframe,
discardable,
payload,
},
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __pauseOpenRtcMediaPublication)]
pub fn pause_openrtc_media_publication(
&self,
publication_id: String,
) -> Result<(), JsValue> {
let publication_id: crate::media::PublicationId =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.portable_media
.borrow_mut()
.pause_publication(publication_id)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __retireOpenRtcMediaPublication)]
pub fn retire_openrtc_media_publication(
&self,
publication_id: String,
) -> Result<(), JsValue> {
let publication_id: crate::media::PublicationId =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.portable_media
.borrow_mut()
.retire_publication(publication_id);
Ok(())
}
#[wasm_bindgen(js_name = __retireOpenRtcMediaReceiver)]
pub fn retire_openrtc_media_receiver(
&self,
publication_id: String,
media_generation: u32,
) -> Result<(), JsValue> {
let publication_id: crate::media::PublicationId =
publication_id
.parse()
.map_err(|error: crate::media::MediaProtocolError| {
JsValue::from_str(&error.to_string())
})?;
self.portable_media
.borrow_mut()
.retire_receiver(publication_id, media_generation);
Ok(())
}
#[wasm_bindgen(js_name = __decodeOpenRtcMediaChunk)]
pub fn decode_openrtc_media_chunk(&self, encoded: Vec<u8>) -> Result<JsValue, JsValue> {
let chunk = self
.portable_media
.borrow_mut()
.decode_chunk(&encoded)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
crate::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
.and_then(|_| {
crate::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
})
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let object = js_sys::Object::new();
let set = |name: &str, value: JsValue| -> Result<(), JsValue> {
js_sys::Reflect::set(&object, &JsValue::from_str(name), &value).map(|_| ())
};
set(
"publicationId",
JsValue::from_str(&chunk.publication_id.to_string()),
)?;
set(
"mediaGeneration",
JsValue::from_f64(chunk.media_generation as f64),
)?;
set("sequence", JsValue::from_f64(chunk.sequence as f64))?;
set("timestampUs", JsValue::from_f64(chunk.timestamp_us as f64))?;
set("durationUs", JsValue::from_f64(chunk.duration_us as f64))?;
set(
"kind",
JsValue::from_str(match chunk.kind {
crate::media::MediaKind::Audio => "audio",
crate::media::MediaKind::Video => "video",
crate::media::MediaKind::Screen => "screen",
crate::media::MediaKind::Data => "data",
}),
)?;
set(
"codec",
JsValue::from_str(match chunk.codec {
crate::media::MediaCodec::Opus => "opus",
crate::media::MediaCodec::H264 => "h264",
crate::media::MediaCodec::Vp8 => "vp8",
crate::media::MediaCodec::Vp9 => "vp9",
crate::media::MediaCodec::Av1 => "av1",
crate::media::MediaCodec::Pcm => "pcm",
crate::media::MediaCodec::Opaque => "opaque",
}),
)?;
set("keyframe", JsValue::from_bool(chunk.keyframe))?;
set("discardable", JsValue::from_bool(chunk.discardable))?;
set(
"payload",
js_sys::Uint8Array::from(chunk.payload.as_slice()).into(),
)?;
Ok(object.into())
}
#[wasm_bindgen(js_name = __decodeOpenRtcMediaControl)]
pub fn decode_openrtc_media_control(&self, encoded: Vec<u8>) -> Result<JsValue, JsValue> {
let control = self
.portable_media
.borrow_mut()
.decode_control(&encoded)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let portable = match control {
crate::media::MediaControlFrame::Publish {
publication_id,
media_generation,
kind,
codec,
clock_rate,
coded_width,
coded_height,
channels,
} => serde_json::json!({
"type": "publish",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"kind": kind,
"codec": codec,
"clockRate": clock_rate,
"codedWidth": coded_width,
"codedHeight": coded_height,
"channels": channels,
}),
crate::media::MediaControlFrame::SetEnabled {
publication_id,
media_generation,
enabled,
} => serde_json::json!({
"type": "set-enabled",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"enabled": enabled,
}),
crate::media::MediaControlFrame::RequestKeyframe {
publication_id,
media_generation,
} => serde_json::json!({
"type": "request-keyframe",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
}),
crate::media::MediaControlFrame::Stop {
publication_id,
media_generation,
reason,
} => serde_json::json!({
"type": "stop",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"reason": reason,
}),
};
serde::Serialize::serialize(
&portable,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __setIrohCarrierActionHandler)]
pub fn set_iroh_carrier_action_handler(&self, handler: Option<js_sys::Function>) {
*self.wasm_carrier_action_handler.borrow_mut() = handler;
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __getIrohCarrierReadiness)]
pub async fn get_iroh_carrier_readiness(&self) -> Result<JsValue, JsValue> {
let readiness = serde_json::json!({
"webrtcConfigured": self.inner.is_webrtc_carrier_configured().await,
"moqConfigured": self.inner.is_moq_carrier_configured().await,
"webrtcAvailable": self.inner.is_webrtc_carrier_enabled().await,
"moqAvailable": self.inner.is_moq_carrier_enabled().await,
"actionHandlerInstalled": self.wasm_carrier_action_handler.borrow().is_some(),
});
serde::Serialize::serialize(
&readiness,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __configureWebRtcCarrier)]
pub async fn configure_iroh_webrtc_carrier(
&self,
enabled: bool,
privacy_mode: bool,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.inner.set_webrtc_carrier(enabled, privacy_mode).await;
Ok(())
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __configureIrohRoutePolicy)]
pub async fn configure_iroh_route_policy(
&self,
relay_only: bool,
optimize_for: Option<String>,
route_priority: JsValue,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let optimize_for = match optimize_for.as_deref().unwrap_or("balanced") {
"balanced" => crate::route_policy::RouteOptimization::Balanced,
"lowest-latency" => crate::route_policy::RouteOptimization::LowestLatency,
_ => {
return Err(JsValue::from_str(
"transport optimization must be balanced or lowest-latency",
))
}
};
let route_priority = if route_priority.is_null() || route_priority.is_undefined() {
Vec::new()
} else {
serde_wasm_bindgen::from_value(route_priority)
.map_err(|error| JsValue::from_str(&error.to_string()))?
};
self.inner
.configure_wasm_route_policy(relay_only, optimize_for, route_priority)
.await;
Ok(())
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __applyIrohCarrierPolicy)]
pub async fn apply_iroh_carrier_policy(&self) -> Result<u32, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let capabilities: Vec<_> = self
.wasm_remote_carrier_capabilities
.borrow()
.iter()
.map(|(connection_id, capabilities)| (connection_id.clone(), *capabilities))
.collect();
let mut reconciled = 0u32;
for (connection_id, capabilities) in capabilities {
let Some(remote_endpoint_id) = self
.inner
.wasm_iroh_carrier_remote_endpoint_id(&connection_id)
.await
else {
continue;
};
if self
.begin_preferred_iroh_carrier_locked(
connection_id,
remote_endpoint_id,
capabilities.webrtc,
capabilities.moq,
)
.await?
{
reconciled = reconciled.saturating_add(1);
}
}
Ok(reconciled)
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __reconcileIrohCarriersAfterNetworkChange)]
pub async fn reconcile_iroh_carriers_after_network_change(&self) -> u32 {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let mut reconciled = 0u32;
let mut recoveries = Vec::new();
#[cfg(feature = "transport-webrtc")]
{
let attempts: Vec<_> = self
.wasm_webrtc_carrier_attempts
.borrow()
.values()
.cloned()
.collect();
let probed = futures::future::join_all(attempts.into_iter().map(|attempt| {
let inner = self.inner.clone();
async move {
let needs_recovery = inner
.wasm_carrier_generation_needs_network_recovery(
&attempt.connection_id,
attempt.remote_endpoint_id,
crate::client::IrohPathKind::WebRtc,
attempt.generation.transport_generation.saturating_add(1),
)
.await;
(attempt, needs_recovery)
}
}))
.await;
for (attempt, needs_recovery) in probed {
if !needs_recovery {
continue;
}
if retire_selected_wasm_webrtc_carrier(
self.inner.clone(),
self.wasm_webrtc_carrier_attempts.clone(),
self.wasm_webrtc_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
&attempt,
"network-change",
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await
{
reconciled = reconciled.saturating_add(1);
recoveries.push((
attempt.connection_id.clone(),
attempt.remote_endpoint_id,
attempt.generation.transport_generation.saturating_add(1),
));
}
}
}
#[cfg(feature = "transport-moq")]
{
let attempts: Vec<_> = self
.wasm_moq_carrier_attempts
.borrow()
.values()
.cloned()
.collect();
let probed = futures::future::join_all(attempts.into_iter().map(|attempt| {
let inner = self.inner.clone();
async move {
let needs_recovery = inner
.wasm_carrier_generation_needs_network_recovery(
&attempt.connection_id,
attempt.remote_endpoint_id,
crate::client::IrohPathKind::Moq,
attempt.generation.transport_generation.saturating_add(1),
)
.await;
(attempt, needs_recovery)
}
}))
.await;
for (attempt, needs_recovery) in probed {
if !needs_recovery {
continue;
}
if retire_selected_wasm_moq_carrier(
self.inner.clone(),
self.wasm_moq_carrier_attempts.clone(),
self.wasm_moq_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
&attempt,
"network-change",
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await
{
reconciled = reconciled.saturating_add(1);
recoveries.push((
attempt.connection_id.clone(),
attempt.remote_endpoint_id,
attempt.generation.transport_generation.saturating_add(1),
));
}
}
}
for (connection_id, remote_endpoint_id, transport_generation) in recoveries {
let recovered = self
.inner
.recover_retired_wasm_carrier_generation(
&connection_id,
remote_endpoint_id,
transport_generation,
)
.await;
if !recovered {
web_sys::console::warn_1(&JsValue::from_str(&format!(
"[OpenRTC][Iroh carrier] base recovery remained pending connection_id={connection_id} transport_generation={transport_generation}"
)));
}
}
if reconciled > 0 {
self.inner.wake_browser_auto_connect();
}
reconciled
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __forgetIrohCarrierPeer)]
pub async fn forget_iroh_carrier_peer(&self, connection_id: String) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.forget_iroh_carrier_peer_locked(&connection_id);
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
fn forget_iroh_carrier_peer_locked(&self, connection_id: &str) {
self.wasm_remote_carrier_capabilities
.borrow_mut()
.remove(connection_id);
self.inner
.forget_iroh_carrier_peer_policy_epoch(connection_id);
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __retireIrohCarrierPeer)]
pub async fn retire_iroh_carrier_peer(
&self,
connection_id: String,
terminal_reason: Option<String>,
) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
#[cfg(feature = "transport-webrtc")]
self.retire_iroh_webrtc_carrier_locked(
connection_id.clone(),
terminal_reason.as_deref(),
)
.await;
#[cfg(feature = "transport-moq")]
self.retire_iroh_moq_carrier_locked(connection_id.clone(), terminal_reason.as_deref())
.await;
self.forget_iroh_carrier_peer_locked(&connection_id);
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
fn has_matching_in_flight_wasm_carrier_attempt(
&self,
connection_id: &str,
capabilities: WasmRemoteCarrierCapabilities,
) -> bool {
#[cfg(feature = "transport-webrtc")]
if capabilities.webrtc
&& self
.wasm_webrtc_carrier_attempts
.borrow()
.get(connection_id)
.is_some_and(|attempt| {
attempt.role == "responder" && !attempt.completion_started
})
{
return true;
}
#[cfg(feature = "transport-moq")]
if capabilities.moq
&& self
.wasm_moq_carrier_attempts
.borrow()
.get(connection_id)
.is_some_and(|attempt| attempt.role == "responder")
{
return true;
}
false
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __beginPreferredIrohCarrier)]
pub async fn begin_preferred_iroh_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
remote_supports_webrtc: bool,
remote_supports_moq: bool,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.begin_preferred_iroh_carrier_locked(
connection_id,
remote_endpoint_id,
remote_supports_webrtc,
remote_supports_moq,
)
.await
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
async fn begin_preferred_iroh_carrier_locked(
&self,
connection_id: String,
remote_endpoint_id: String,
remote_supports_webrtc: bool,
remote_supports_moq: bool,
) -> Result<bool, JsValue> {
let capabilities = WasmRemoteCarrierCapabilities {
webrtc: remote_supports_webrtc,
moq: remote_supports_moq,
};
let previous_capabilities = self
.wasm_remote_carrier_capabilities
.borrow()
.get(&connection_id)
.copied();
if previous_capabilities != Some(capabilities) {
let matching_first_attempt = previous_capabilities.is_none()
&& self
.has_matching_in_flight_wasm_carrier_attempt(&connection_id, capabilities);
if crate::client::should_bump_remote_carrier_peer_policy_epoch(
previous_capabilities.is_some(),
true,
matching_first_attempt,
) {
self.inner
.bump_wasm_iroh_carrier_peer_policy_epoch(&connection_id)
.await;
}
self.wasm_remote_carrier_capabilities
.borrow_mut()
.insert(connection_id.clone(), capabilities);
}
if self
.inner
.reconcile_wasm_iroh_carrier_policy(
&connection_id,
remote_supports_webrtc,
remote_supports_moq,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?
{
return Ok(true);
}
let ranked = self
.inner
.ranked_wasm_iroh_carriers(remote_supports_webrtc, remote_supports_moq)
.await;
match ranked.first() {
#[cfg(feature = "transport-webrtc")]
Some(crate::route_policy::KnownRoute::WebRtc) => {
self.begin_iroh_webrtc_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
#[cfg(feature = "transport-moq")]
Some(crate::route_policy::KnownRoute::Moq) => {
self.begin_iroh_moq_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
_ => Ok(false),
}
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
#[wasm_bindgen(js_name = __advancePreferredIrohCarrier)]
pub async fn advance_preferred_iroh_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
failed_route: String,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let Some(capabilities) = self
.wasm_remote_carrier_capabilities
.borrow()
.get(&connection_id)
.copied()
else {
return Ok(false);
};
let failed = crate::route_policy::normalize_route(&failed_route)
.ok_or_else(|| JsValue::from_str("unknown failed carrier route"))?;
let ranked = self
.inner
.ranked_wasm_iroh_carriers(capabilities.webrtc, capabilities.moq)
.await;
let Some(next_index) = ranked
.iter()
.position(|route| *route == failed)
.map(|index| index + 1)
.filter(|index| *index < ranked.len())
else {
return Ok(false);
};
let next = ranked[next_index];
match next {
#[cfg(feature = "transport-webrtc")]
crate::route_policy::KnownRoute::WebRtc => {
self.begin_iroh_webrtc_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
#[cfg(feature = "transport-moq")]
crate::route_policy::KnownRoute::Moq => {
self.begin_iroh_moq_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
_ => Ok(false),
}
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __configureMoqCarrier)]
pub async fn configure_iroh_moq_carrier(
&self,
enabled: bool,
relay_url: Option<String>,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let relay_url = relay_url.map(|value| value.trim().to_string());
if enabled && relay_url.as_deref().unwrap_or_default().is_empty() {
return Err(JsValue::from_str(
"MoQ Iroh carrier requires an explicit relay URL",
));
}
self.inner.set_moq_carrier(enabled, relay_url).await;
Ok(())
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __beginMoqCarrier)]
pub async fn begin_iroh_moq_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.begin_iroh_moq_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
#[cfg(feature = "transport-moq")]
async fn begin_iroh_moq_carrier_attempt(
&self,
connection_id: String,
remote_endpoint_id: String,
retry_count: u8,
) -> Result<bool, JsValue> {
if !self.inner.is_moq_carrier_enabled().await {
return Ok(false);
}
let local_endpoint_id = self
.inner
.current_node_id()
.await
.ok_or_else(|| JsValue::from_str("local Iroh endpoint id is unavailable"))?;
if local_endpoint_id.as_str() <= remote_endpoint_id.as_str()
|| !crate::iroh_connection_policy::custom_carrier_base_allows(
self.inner.iroh_path_kind(&remote_endpoint_id).await,
crate::client::IrohPathKind::Moq,
false,
)
{
return Ok(false);
}
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.inner
.get_connection(endpoint_id)
.await
.ok_or_else(|| JsValue::from_str("MoQ carrier base is unavailable"))?;
let generation = self
.inner
.current_wasm_peer_data_generation(&connection_id, None)
.await
.ok_or_else(|| JsValue::from_str("MoQ carrier generation is unavailable"))?;
let bootstrap = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::request(
crate::iroh_carrier_bootstrap::CarrierBootstrapKind::MoqDraft14,
crate::iroh_carrier_bootstrap::CarrierGenerationFence {
transport_stable_id: generation.transport_stable_id,
transport_generation: generation.transport_generation,
route_generation: generation.route_generation,
},
retry_count.saturating_add(1),
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let kind = crate::client::IrohPathKind::Moq;
if !matches!(
self.inner
.reserve_wasm_carrier_upgrade(
&connection_id,
kind,
&bootstrap.upgrade_id,
generation,
)
.await,
crate::client::WasmCarrierUpgradeReservation::Reserved
) {
let pending = self
.wasm_moq_carrier_attempts
.borrow_mut()
.get_mut(&connection_id)
.filter(|attempt| {
attempt.role == "initiator"
&& attempt.generation == generation
&& attempt.prepared
&& !attempt.retry_sent
})
.map(|attempt| {
attempt.retry_sent = true;
attempt.clone()
});
if let Some(pending) = pending {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": pending.remote_endpoint_id.to_string(),
"envelope": pending.bootstrap,
}),
);
return Ok(true);
}
return Ok(false);
}
let previous = self.wasm_moq_carrier_attempts.borrow_mut().insert(
connection_id.clone(),
WasmMoqCarrierAttempt {
connection_id: connection_id.clone(),
remote_endpoint_id: endpoint_id,
bootstrap: bootstrap.clone(),
generation,
role: "initiator",
prepared: false,
retry_sent: false,
retry_count,
inbound_authorization_expires_at_ms: None,
},
);
if let Some(previous) = previous {
self.wasm_moq_carrier_sessions
.borrow_mut()
.remove(&previous.bootstrap.upgrade_id);
}
let (publish_namespace, subscribe_namespace, track_name) = wasm_moq_carrier_namespaces(
&local_endpoint_id,
&remote_endpoint_id,
&bootstrap.carrier_session_id,
);
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "prepare-moq",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"role": "initiator",
"upgradeId": bootstrap.upgrade_id,
"carrierSessionId": bootstrap.carrier_session_id,
"transportGeneration": generation.transport_generation.saturating_add(1),
"publishNamespace": publish_namespace,
"subscribeNamespace": subscribe_namespace,
"trackName": track_name,
}),
);
self.schedule_wasm_moq_carrier_watchdog(bootstrap);
Ok(true)
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __moqCarrierPrepared)]
pub async fn iroh_moq_carrier_prepared(
&self,
connection_id: String,
upgrade_id: String,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let attempt = {
let mut attempts = self.wasm_moq_carrier_attempts.borrow_mut();
let attempt = attempts
.get_mut(&connection_id)
.filter(|attempt| attempt.bootstrap.upgrade_id == upgrade_id)
.ok_or_else(|| JsValue::from_str("MoQ carrier attempt is stale"))?;
attempt.prepared = true;
attempt.clone()
};
if !self
.inner
.wasm_carrier_upgrade_is_current(
&connection_id,
crate::client::IrohPathKind::Moq,
&upgrade_id,
attempt.generation,
)
.await
{
return Err(JsValue::from_str("MoQ carrier attempt was retired"));
}
if attempt.role == "initiator" {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": attempt.bootstrap,
}),
);
} else {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "activate-moq",
"connectionId": attempt.connection_id,
"upgradeId": upgrade_id,
}),
);
}
Ok(())
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __moqCarrierFailed)]
pub async fn iroh_moq_carrier_failed(
&self,
connection_id: String,
upgrade_id: String,
failure_code: String,
) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let attempt = self
.wasm_moq_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| attempt.bootstrap.upgrade_id == upgrade_id)
.cloned();
if let Some(attempt) = attempt {
let failure_code = match failure_code.as_str() {
"relay-failed" => "relay-failed",
"draft14-datagram-unavailable" => "draft14-datagram-unavailable",
"carrier-backpressure" => "carrier-backpressure",
_ => "browser-adapter-failed",
};
if retire_selected_wasm_moq_carrier(
self.inner.clone(),
self.wasm_moq_carrier_attempts.clone(),
self.wasm_moq_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
&attempt,
failure_code,
crate::lifecycle_reason::REASON_IROH_CARRIER_FAILED,
)
.await
{
return;
}
self.fail_wasm_moq_carrier_attempt(attempt, failure_code, true)
.await;
}
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __retireMoqCarrier)]
pub async fn retire_iroh_moq_carrier(
&self,
connection_id: String,
terminal_reason: Option<String>,
) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.retire_iroh_moq_carrier_locked(connection_id, terminal_reason.as_deref())
.await;
}
#[cfg(feature = "transport-moq")]
async fn retire_iroh_moq_carrier_locked(
&self,
connection_id: String,
terminal_reason: Option<&str>,
) {
let attempt = self
.wasm_moq_carrier_attempts
.borrow_mut()
.remove(&connection_id);
if let Some(attempt) = attempt {
self.inner
.retire_wasm_carrier_upgrade(
&connection_id,
crate::client::IrohPathKind::Moq,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await;
let carrier = self
.wasm_moq_carrier_sessions
.borrow_mut()
.remove(&attempt.bootstrap.upgrade_id);
if let (Some(reason), Some(session)) = (terminal_reason, carrier.as_ref()) {
let _ = session.send_terminal(reason).await;
}
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "retire-moq",
"connectionId": connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": "logical-connection-retired",
}),
);
}
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __beginWebRtcCarrier)]
pub async fn begin_iroh_webrtc_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.begin_iroh_webrtc_carrier_attempt(connection_id, remote_endpoint_id, 0)
.await
}
#[cfg(feature = "transport-webrtc")]
async fn begin_iroh_webrtc_carrier_attempt(
&self,
connection_id: String,
remote_endpoint_id: String,
retry_count: u8,
) -> Result<bool, JsValue> {
if !self.inner.is_webrtc_carrier_enabled().await {
return Ok(false);
}
let local_endpoint_id = self
.inner
.current_node_id()
.await
.ok_or_else(|| JsValue::from_str("local Iroh endpoint id is unavailable"))?;
if local_endpoint_id.as_str() <= remote_endpoint_id.as_str()
|| !crate::iroh_connection_policy::custom_carrier_base_allows(
self.inner.iroh_path_kind(&remote_endpoint_id).await,
crate::client::IrohPathKind::WebRtc,
false,
)
{
return Ok(false);
}
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.inner
.get_connection(endpoint_id)
.await
.ok_or_else(|| JsValue::from_str("WebRTC carrier base is unavailable"))?;
let generation = self
.inner
.current_wasm_peer_data_generation(&connection_id, None)
.await
.ok_or_else(|| JsValue::from_str("WebRTC carrier generation is unavailable"))?;
let bootstrap = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::request(
crate::iroh_carrier_bootstrap::CarrierBootstrapKind::WebRtc,
crate::iroh_carrier_bootstrap::CarrierGenerationFence {
transport_stable_id: generation.transport_stable_id,
transport_generation: generation.transport_generation,
route_generation: generation.route_generation,
},
1,
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let kind = crate::client::IrohPathKind::WebRtc;
if !matches!(
self.inner
.reserve_wasm_carrier_upgrade(
&connection_id,
kind,
&bootstrap.upgrade_id,
generation,
)
.await,
crate::client::WasmCarrierUpgradeReservation::Reserved
) {
let pending = self
.wasm_webrtc_carrier_attempts
.borrow_mut()
.get_mut(&connection_id)
.filter(|attempt| {
attempt.role == "initiator"
&& attempt.generation == generation
&& attempt.prepared
&& !attempt.retry_sent
&& !attempt.completion_started
})
.map(|attempt| {
attempt.retry_sent = true;
attempt.clone()
});
if let Some(pending) = pending {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": pending.remote_endpoint_id.to_string(),
"envelope": pending.bootstrap,
}),
);
return Ok(true);
}
return Ok(false);
}
let previous = self.wasm_webrtc_carrier_attempts.borrow_mut().insert(
connection_id.clone(),
WasmWebRtcCarrierAttempt {
connection_id: connection_id.clone(),
remote_endpoint_id: endpoint_id,
bootstrap: bootstrap.clone(),
generation,
role: "initiator",
prepared: false,
retry_sent: false,
offer_started: false,
remote_ready: false,
completion_started: false,
retry_count,
inbound_authorization_expires_at_ms: None,
},
);
if let Some(previous) = previous {
self.wasm_webrtc_carrier_sessions
.borrow_mut()
.remove(&previous.bootstrap.upgrade_id);
}
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "prepare-webrtc",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"role": "initiator",
"upgradeId": bootstrap.upgrade_id,
"carrierSessionId": bootstrap.carrier_session_id,
"transportGeneration": generation.transport_generation.saturating_add(1),
}),
);
self.schedule_wasm_webrtc_carrier_watchdog(bootstrap);
Ok(true)
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __webRtcCarrierPrepared)]
pub async fn iroh_webrtc_carrier_prepared(
&self,
connection_id: String,
upgrade_id: String,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let attempt = {
let mut attempts = self.wasm_webrtc_carrier_attempts.borrow_mut();
let attempt = attempts
.get_mut(&connection_id)
.filter(|attempt| attempt.bootstrap.upgrade_id == upgrade_id)
.ok_or_else(|| JsValue::from_str("WebRTC carrier attempt is stale"))?;
attempt.prepared = true;
attempt.clone()
};
let kind = crate::client::IrohPathKind::WebRtc;
if !self
.inner
.wasm_carrier_upgrade_is_current(
&connection_id,
kind,
&upgrade_id,
attempt.generation,
)
.await
{
return Err(JsValue::from_str("WebRTC carrier attempt was retired"));
}
if attempt.role == "initiator" {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": attempt.bootstrap,
}),
);
} else {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": attempt.connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": {
"type": "#pluto-signal",
"content": {
"transport": "webrtc",
"type": "renegotiate",
"negotiationId": upgrade_id,
}
},
}),
);
}
Ok(())
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __webRtcCarrierFailed)]
pub async fn iroh_webrtc_carrier_failed(
&self,
connection_id: String,
upgrade_id: String,
failure_code: String,
) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let attempt = self
.wasm_webrtc_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| attempt.bootstrap.upgrade_id == upgrade_id)
.cloned();
if let Some(attempt) = attempt {
let failure_code = match failure_code.as_str() {
"data-channel-failed" => "data-channel-failed",
"ice-failed" => "ice-failed",
"signaling-failed" => "signaling-failed",
"carrier-backpressure" => "carrier-backpressure",
"carrier-base-generation-stale" => "carrier-base-generation-stale",
_ => "browser-adapter-failed",
};
if retire_selected_wasm_webrtc_carrier(
self.inner.clone(),
self.wasm_webrtc_carrier_attempts.clone(),
self.wasm_webrtc_carrier_sessions.clone(),
self.wasm_carrier_action_handler.clone(),
&attempt,
failure_code,
crate::lifecycle_reason::REASON_IROH_CARRIER_FAILED,
)
.await
{
return;
}
let should_retry = attempt.role == "initiator"
&& attempt.retry_count == 0
&& matches!(
failure_code,
"data-channel-failed"
| "ice-failed"
| "signaling-failed"
| "carrier-base-generation-stale"
);
let remote_endpoint_id = attempt.remote_endpoint_id.to_string();
self.fail_wasm_webrtc_carrier_attempt(attempt, failure_code, true)
.await;
if should_retry {
gloo_timers::future::sleep(std::time::Duration::from_millis(500)).await;
let _ = self
.begin_iroh_webrtc_carrier_attempt(connection_id, remote_endpoint_id, 1)
.await;
}
}
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __retireWebRtcCarrier)]
pub async fn retire_iroh_webrtc_carrier(
&self,
connection_id: String,
terminal_reason: Option<String>,
) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
self.retire_iroh_webrtc_carrier_locked(connection_id, terminal_reason.as_deref())
.await;
}
#[cfg(feature = "transport-webrtc")]
async fn retire_iroh_webrtc_carrier_locked(
&self,
connection_id: String,
terminal_reason: Option<&str>,
) {
let attempt = self
.wasm_webrtc_carrier_attempts
.borrow_mut()
.remove(&connection_id);
if let Some(attempt) = attempt {
self.inner
.retire_wasm_carrier_upgrade(
&connection_id,
crate::client::IrohPathKind::WebRtc,
&attempt.bootstrap.upgrade_id,
attempt.generation,
)
.await;
let carrier = self
.wasm_webrtc_carrier_sessions
.borrow_mut()
.remove(&attempt.bootstrap.upgrade_id);
if let (Some(reason), Some(session)) = (terminal_reason, carrier.as_ref()) {
let _ = session.send_terminal(reason).await;
}
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "retire-webrtc",
"connectionId": connection_id,
"upgradeId": attempt.bootstrap.upgrade_id,
"failureCode": "logical-connection-retired",
}),
);
}
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __handleIrohCarrierControl)]
pub async fn handle_iroh_carrier_control(
&self,
connection_id: String,
remote_endpoint_id: String,
frame: JsValue,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let frame: serde_json::Value = serde_wasm_bindgen::from_value(frame)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
if frame.get("type").and_then(serde_json::Value::as_str) == Some("#pluto-signal")
&& frame
.get("content")
.and_then(|content| content.get("transport"))
.and_then(serde_json::Value::as_str)
== Some("webrtc")
{
let negotiation_id = frame
.get("content")
.and_then(|content| content.get("negotiationId"))
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
let current = self
.wasm_webrtc_carrier_attempts
.borrow()
.get(&connection_id)
.is_some_and(|attempt| {
attempt.bootstrap.upgrade_id == negotiation_id
&& attempt.remote_endpoint_id.to_string() == remote_endpoint_id
});
if current {
let signal_type = frame
.get("content")
.and_then(|content| content.get("type"))
.and_then(serde_json::Value::as_str);
let start_offer = if signal_type == Some("renegotiate") {
self.wasm_webrtc_carrier_attempts
.borrow_mut()
.get_mut(&connection_id)
.filter(|attempt| {
attempt.bootstrap.upgrade_id == negotiation_id
&& attempt.role == "initiator"
&& !attempt.offer_started
})
.map(|attempt| {
attempt.offer_started = true;
})
.is_some()
} else {
false
};
if start_offer {
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "start-webrtc",
"connectionId": connection_id,
"upgradeId": negotiation_id,
}),
);
return Ok(true);
}
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "apply-webrtc-signal",
"connectionId": connection_id,
"upgradeId": negotiation_id,
"signal": frame,
}),
);
}
return Ok(true);
}
let Some(bootstrap) =
crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::from_json(&frame)
else {
return Ok(false);
};
#[cfg(feature = "transport-moq")]
if bootstrap.carrier == crate::iroh_carrier_bootstrap::CarrierBootstrapKind::MoqDraft14
{
return self
.handle_wasm_moq_bootstrap(connection_id, remote_endpoint_id, bootstrap)
.await;
}
if bootstrap.carrier != crate::iroh_carrier_bootstrap::CarrierBootstrapKind::WebRtc {
return Ok(false);
}
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let kind = crate::client::IrohPathKind::WebRtc;
match bootstrap.action {
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Request => {
let local_endpoint_id = self
.inner
.current_node_id()
.await
.ok_or_else(|| JsValue::from_str("local endpoint id is unavailable"))?;
if local_endpoint_id.as_str() >= remote_endpoint_id.as_str()
|| !self.inner.is_webrtc_carrier_enabled().await
{
return Ok(true);
}
if !crate::iroh_connection_policy::custom_carrier_base_allows(
self.inner.iroh_path_kind(&remote_endpoint_id).await,
crate::client::IrohPathKind::WebRtc,
false,
) {
return Ok(true);
}
let base_connection = self
.inner
.get_connection(endpoint_id)
.await
.ok_or_else(|| JsValue::from_str("WebRTC carrier base is unavailable"))?;
let generation = self
.inner
.current_wasm_peer_data_generation(&connection_id, None)
.await
.ok_or_else(|| {
JsValue::from_str("WebRTC carrier generation is unavailable")
})?;
if crate::transport_generation::for_connection(&base_connection)
!= generation.transport_stable_id
{
let failed =
crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::failed_from(
&bootstrap,
"carrier-base-generation-stale",
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"envelope": failed,
}),
);
return Ok(true);
}
if !matches!(
self.inner
.reserve_wasm_carrier_upgrade(
&connection_id,
kind,
&bootstrap.upgrade_id,
generation,
)
.await,
crate::client::WasmCarrierUpgradeReservation::Reserved
) {
return Ok(true);
}
let node = self
.inner
.iroh_node
.read()
.await
.as_ref()
.cloned()
.ok_or_else(|| JsValue::from_str("Iroh node is unavailable"))?;
let authorization_expiry = node
.authorize_pending_inbound_replacement(
endpoint_id,
crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID,
std::time::Duration::from_secs(45),
)
.await;
let previous = self.wasm_webrtc_carrier_attempts.borrow_mut().insert(
connection_id.clone(),
WasmWebRtcCarrierAttempt {
connection_id: connection_id.clone(),
remote_endpoint_id: endpoint_id,
bootstrap: bootstrap.clone(),
generation,
role: "responder",
prepared: false,
retry_sent: false,
offer_started: false,
remote_ready: false,
completion_started: false,
retry_count: bootstrap.attempt.saturating_sub(1),
inbound_authorization_expires_at_ms: Some(authorization_expiry),
},
);
if let Some(previous) = previous {
self.wasm_webrtc_carrier_sessions
.borrow_mut()
.remove(&previous.bootstrap.upgrade_id);
}
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "prepare-webrtc",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"role": "responder",
"upgradeId": bootstrap.upgrade_id,
"carrierSessionId": bootstrap.carrier_session_id,
"transportGeneration": generation.transport_generation.saturating_add(1),
}),
);
self.schedule_wasm_webrtc_carrier_watchdog(bootstrap);
}
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Ready => {
let upgrade_id = {
let mut attempts = self.wasm_webrtc_carrier_attempts.borrow_mut();
attempts
.get_mut(&connection_id)
.filter(|attempt| bootstrap.is_response_to(&attempt.bootstrap))
.map(|attempt| {
attempt.remote_ready = true;
attempt.bootstrap.upgrade_id.clone()
})
};
let attempt = upgrade_id.as_deref().and_then(|upgrade_id| {
self.take_ready_outbound_wasm_webrtc_carrier_attempt(
&connection_id,
upgrade_id,
)
});
if let Some(attempt) = attempt {
self.spawn_outbound_wasm_webrtc_carrier_completion(attempt);
}
}
crate::iroh_carrier_bootstrap::CarrierBootstrapAction::Failed => {
let attempt = self
.wasm_webrtc_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| bootstrap.is_response_to(&attempt.bootstrap))
.cloned();
if let Some(attempt) = attempt {
let failure_code = crate::client::wasm_peer_carrier_failure_code(
bootstrap.failure_code.as_deref(),
);
if crate::client::should_rearm_wasm_carrier_event_retry(
failure_code,
attempt.role == "initiator",
attempt.completion_started,
) {
if let Some(current) = self
.wasm_webrtc_carrier_attempts
.borrow_mut()
.get_mut(&connection_id)
.filter(|current| {
current.bootstrap.upgrade_id == attempt.bootstrap.upgrade_id
&& !current.completion_started
})
{
current.retry_sent = false;
}
return Ok(true);
}
let should_retry = attempt.role == "initiator"
&& attempt.retry_count == 0
&& failure_code == "carrier-base-generation-stale";
let remote_endpoint_id = attempt.remote_endpoint_id.to_string();
self.fail_wasm_webrtc_carrier_attempt(attempt, failure_code, false)
.await;
if should_retry {
gloo_timers::future::sleep(std::time::Duration::from_millis(500)).await;
let _ = self
.begin_iroh_webrtc_carrier_attempt(
connection_id,
remote_endpoint_id,
1,
)
.await;
}
}
}
}
Ok(true)
}
#[cfg(all(feature = "transport-moq", not(feature = "transport-webrtc")))]
#[wasm_bindgen(js_name = __handleIrohCarrierControl)]
pub async fn handle_iroh_carrier_control_moq_only(
&self,
connection_id: String,
remote_endpoint_id: String,
frame: JsValue,
) -> Result<bool, JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let frame: serde_json::Value = serde_wasm_bindgen::from_value(frame)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let Some(bootstrap) =
crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::from_json(&frame)
else {
return Ok(false);
};
if bootstrap.carrier != crate::iroh_carrier_bootstrap::CarrierBootstrapKind::MoqDraft14
{
return Ok(false);
}
self.handle_wasm_moq_bootstrap(connection_id, remote_endpoint_id, bootstrap)
.await
}
#[cfg(feature = "transport-webrtc")]
#[wasm_bindgen(js_name = __attachWebRtcCarrier)]
pub async fn attach_iroh_webrtc_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
channel: web_sys::RtcDataChannel,
upgrade_id: String,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let attempt = self
.wasm_webrtc_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| {
attempt.bootstrap.upgrade_id == upgrade_id
&& attempt.remote_endpoint_id == endpoint_id
})
.cloned()
.ok_or_else(|| JsValue::from_str("WebRTC carrier attempt is stale"))?;
if !self
.inner
.wasm_carrier_upgrade_is_current(
&connection_id,
crate::client::IrohPathKind::WebRtc,
&upgrade_id,
attempt.generation,
)
.await
{
return Err(JsValue::from_str("WebRTC carrier attempt was retired"));
}
let packet_session = self
.inner
.activate_iroh_packet_carrier(
crate::iroh_carrier_kind::IrohCarrierKind::WebRtc,
endpoint_id,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let expected = attempt
.bootstrap
.frame_expectation()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let application_key = self
.inner
.application_crypto_key_for_connection(Some(&connection_id))
.ok_or_else(|| {
JsValue::from_str("WebRTC carrier requires the admitted application crypto key")
})?;
let terminal_client = self.inner.clone();
let terminal_connection_id = connection_id.clone();
let terminal_upgrade_id = upgrade_id.clone();
let terminal_generation = attempt.generation.transport_generation.saturating_add(1);
let terminal_attempts = self.wasm_webrtc_carrier_attempts.clone();
let terminal_sessions = self.wasm_webrtc_carrier_sessions.clone();
let terminal_handler = self.wasm_carrier_action_handler.clone();
let terminal_lifecycle = self.wasm_carrier_peer_lifecycle.clone();
let on_terminal: Rc<dyn Fn(&'static str)> = Rc::new(move |reason| {
let client = terminal_client.clone();
let connection_id = terminal_connection_id.clone();
let upgrade_id = terminal_upgrade_id.clone();
let attempts = terminal_attempts.clone();
let sessions = terminal_sessions.clone();
let handler = terminal_handler.clone();
let lifecycle = terminal_lifecycle.clone();
spawn_local(async move {
let _lifecycle = lifecycle.lock().await;
if !client
.close_current_iroh_carrier_generation_with_reason(
&connection_id,
endpoint_id,
crate::client::IrohPathKind::WebRtc,
terminal_generation,
reason,
)
.await
{
return;
}
if attempts
.borrow()
.get(&connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == upgrade_id)
{
attempts.borrow_mut().remove(&connection_id);
}
sessions.borrow_mut().remove(&upgrade_id);
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-webrtc",
"connectionId": connection_id,
"upgradeId": upgrade_id,
"failureCode": reason,
}),
);
});
});
let carrier = crate::wasm_webrtc_carrier::WasmWebRtcCarrierSession::attach(
channel,
packet_session,
expected,
application_key,
on_terminal,
)?;
self.wasm_webrtc_carrier_sessions
.borrow_mut()
.insert(upgrade_id.clone(), carrier);
if attempt.role == "responder" {
let ready = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::ready_from(
&attempt.bootstrap,
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": remote_endpoint_id,
"envelope": ready,
}),
);
self.spawn_inbound_wasm_webrtc_carrier_completion(attempt);
} else if let Some(attempt) =
self.take_ready_outbound_wasm_webrtc_carrier_attempt(&connection_id, &upgrade_id)
{
self.spawn_outbound_wasm_webrtc_carrier_completion(attempt);
}
Ok(())
}
#[cfg(feature = "transport-moq")]
#[wasm_bindgen(js_name = __attachMoqCarrier)]
pub async fn attach_iroh_moq_carrier(
&self,
connection_id: String,
remote_endpoint_id: String,
datagrams: JsValue,
upgrade_id: String,
) -> Result<(), JsValue> {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let endpoint_id = remote_endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let attempt = self
.wasm_moq_carrier_attempts
.borrow()
.get(&connection_id)
.filter(|attempt| {
attempt.bootstrap.upgrade_id == upgrade_id
&& attempt.remote_endpoint_id == endpoint_id
})
.cloned()
.ok_or_else(|| JsValue::from_str("MoQ carrier attempt is stale"))?;
if !self
.inner
.wasm_carrier_upgrade_is_current(
&connection_id,
crate::client::IrohPathKind::Moq,
&upgrade_id,
attempt.generation,
)
.await
{
return Err(JsValue::from_str("MoQ carrier attempt was retired"));
}
let packet_session = self
.inner
.activate_iroh_packet_carrier(
crate::iroh_carrier_kind::IrohCarrierKind::Moq,
endpoint_id,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let expected = attempt
.bootstrap
.frame_expectation()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let application_key = self
.inner
.application_crypto_key_for_connection(Some(&connection_id))
.ok_or_else(|| {
JsValue::from_str("MoQ carrier requires an installed application key")
})?;
let terminal_client = self.inner.clone();
let terminal_connection_id = connection_id.clone();
let terminal_upgrade_id = upgrade_id.clone();
let terminal_generation = attempt.generation.transport_generation.saturating_add(1);
let terminal_attempts = self.wasm_moq_carrier_attempts.clone();
let terminal_sessions = self.wasm_moq_carrier_sessions.clone();
let terminal_handler = self.wasm_carrier_action_handler.clone();
let terminal_lifecycle = self.wasm_carrier_peer_lifecycle.clone();
let on_terminal: Rc<dyn Fn(&'static str)> = Rc::new(move |reason| {
let client = terminal_client.clone();
let connection_id = terminal_connection_id.clone();
let upgrade_id = terminal_upgrade_id.clone();
let attempts = terminal_attempts.clone();
let sessions = terminal_sessions.clone();
let handler = terminal_handler.clone();
let lifecycle = terminal_lifecycle.clone();
spawn_local(async move {
let _lifecycle = lifecycle.lock().await;
if !client
.close_current_iroh_carrier_generation_with_reason(
&connection_id,
endpoint_id,
crate::client::IrohPathKind::Moq,
terminal_generation,
reason,
)
.await
{
return;
}
if attempts
.borrow()
.get(&connection_id)
.is_some_and(|current| current.bootstrap.upgrade_id == upgrade_id)
{
attempts.borrow_mut().remove(&connection_id);
}
sessions.borrow_mut().remove(&upgrade_id);
emit_wasm_carrier_action(
&handler,
serde_json::json!({
"type": "retire-moq",
"connectionId": connection_id,
"upgradeId": upgrade_id,
"failureCode": reason,
}),
);
});
});
let carrier = crate::wasm_moq_carrier::WasmMoqCarrierSession::attach(
datagrams,
packet_session,
expected,
application_key,
on_terminal,
)?;
if attempt.role == "initiator" {
carrier
.wait_for_peer_data_bidirectional_readiness(std::time::Duration::from_secs(15))
.await?;
}
self.wasm_moq_carrier_sessions
.borrow_mut()
.insert(upgrade_id.clone(), carrier);
if attempt.role == "initiator" {
self.spawn_outbound_wasm_moq_carrier_completion(attempt);
} else {
let ready = crate::iroh_carrier_bootstrap::CarrierBootstrapFrame::ready_from(
&attempt.bootstrap,
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
emit_wasm_carrier_action(
&self.wasm_carrier_action_handler,
serde_json::json!({
"type": "send-control",
"connectionId": connection_id,
"remoteEndpointId": attempt.remote_endpoint_id.to_string(),
"envelope": ready,
}),
);
self.spawn_inbound_wasm_moq_carrier_completion(attempt);
}
Ok(())
}
#[wasm_bindgen(js_name = setIdentityCredential)]
pub fn set_identity_credential(&self, credential: Option<String>) {
let has_credential = credential
.as_ref()
.map(|value| !value.is_empty())
.unwrap_or(false);
let credential_len = credential.as_ref().map(|value| value.len()).unwrap_or(0);
if let Ok(mut guard) = self.identity_credential.lock() {
*guard = credential.filter(|value| !value.is_empty());
}
let should_log = if let Ok(mut guard) = self.last_auth_log.lock() {
let next = (has_credential, credential_len);
if guard.as_ref() == Some(&next) {
false
} else {
*guard = Some(next);
true
}
} else {
true
};
if should_log {
web_sys::console::log_1(&JsValue::from_str(&format!(
"[OPENRTC][WASM-IDENTITY] credential updated present={} len={}",
has_credential, credential_len
)));
}
}
#[wasm_bindgen(js_name = clearIdentityCredential)]
pub fn clear_identity_credential(&self) {
self.set_identity_credential(None);
}
#[wasm_bindgen(js_name = rankRoutes)]
pub fn rank_routes(
&self,
configured_priority: Vec<String>,
candidates: Vec<String>,
) -> Vec<String> {
crate::route_policy::rank_routes(&configured_priority, &candidates)
}
#[wasm_bindgen(js_name = setRelay)]
pub async fn set_relay(&self, enabled: bool) -> Result<(), JsValue> {
self.inner
.set_relay(enabled)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn init_iroh(&self, secret_key: Option<Vec<u8>>) -> Result<String, JsValue> {
let started_at = js_sys::Date::now();
web_sys::console::log_1(&JsValue::from_str(&format!(
"[OPENRTC][WASM-API] init_iroh called has_secret_key={} secret_key_len={}",
secret_key.as_ref().is_some(),
secret_key.as_ref().map(|k| k.len()).unwrap_or(0)
)));
match self.inner.init_iroh(secret_key, vec![]).await {
Ok(node_id) => {
self.inner.clone().start_wasm_accept_bridge();
let elapsed = js_sys::Date::now() - started_at;
web_sys::console::log_1(&JsValue::from_str(&format!(
"[OPENRTC][WASM-API] init_iroh success elapsed_ms={:.0} node_id={}",
elapsed, node_id
)));
Ok(node_id)
}
Err(err) => {
let elapsed = js_sys::Date::now() - started_at;
web_sys::console::error_1(&JsValue::from_str(&format!(
"[OPENRTC][WASM-API] init_iroh failed elapsed_ms={:.0} error={}",
elapsed, err
)));
Err(JsValue::from_str(&err.to_string()))
}
}
}
#[wasm_bindgen(js_name = initIrohWithTestRelay)]
pub async fn init_iroh_with_test_relay(
&self,
secret_key: Option<Vec<u8>>,
test_relay_url: Option<String>,
) -> Result<String, JsValue> {
let node_id = self
.inner
.init_iroh_with_test_relay(secret_key, vec![], test_relay_url.as_deref())
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.inner.clone().start_wasm_accept_bridge();
Ok(node_id)
}
pub async fn iroh_secret_key(&self) -> Result<Vec<u8>, JsValue> {
let node_guard = self.inner.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
Ok(node.secret_key())
} else {
Err(JsValue::from_str("Iroh node not initialized"))
}
}
pub async fn node_addr(&self) -> Result<String, JsValue> {
let node_guard = self.inner.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
let addr = node
.node_addr()
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_json::to_string(&addr).map_err(|e| JsValue::from_str(&e.to_string()))
} else {
Err(JsValue::from_str("Iroh node not initialized"))
}
}
pub async fn endpoint_ticket(&self) -> Result<String, JsValue> {
self.inner
.endpoint_ticket()
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn endpoint_ticket_with_token(
&self,
grant_scope: String,
max_connections: u32,
) -> Result<String, JsValue> {
self.inner
.endpoint_ticket_with_token(&grant_scope, max_connections)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn issue_ticket_invite(
&self,
id: String,
max_peers: u32,
mesh: bool,
) -> Result<String, JsValue> {
let invite = if mesh {
self.inner
.issue_ticket_invite(&id, crate::ticket_mesh::TicketMeshOptions { max_peers })
.await
} else {
self.inner.issue_direct_ticket_invite(&id, max_peers).await
};
invite
.map(|invite| invite.encode())
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn issue_ticket_mesh(
&self,
id: String,
max_peers: u32,
action_handler: js_sys::Function,
) -> Result<String, JsValue> {
if self.ticket_meshes.borrow().contains_key(&id)
|| !self.ticket_mesh_starts.borrow_mut().insert(id.clone())
{
return Err(JsValue::from_str("ticket mesh already has an owner"));
}
let _start = TicketMeshStartGuard {
starts: &self.ticket_mesh_starts,
id: id.clone(),
};
let mesh = crate::ticket_mesh::wasm_session::WasmTicketMesh::issue(
self.inner.clone(),
&id,
crate::ticket_mesh::TicketMeshOptions { max_peers },
action_handler,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let invitation = mesh.invite();
self.ticket_meshes.borrow_mut().insert(id, Rc::new(mesh));
Ok(invitation)
}
pub async fn join_ticket_mesh(
&self,
invitation: String,
action_handler: js_sys::Function,
) -> Result<String, JsValue> {
let invite = crate::ticket_mesh::TicketInvite::parse(&invitation)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let id = invite.id().to_string();
if self.ticket_meshes.borrow().contains_key(&id)
|| !self.ticket_mesh_starts.borrow_mut().insert(id.clone())
{
return Err(JsValue::from_str("ticket mesh already has an owner"));
}
let _start = TicketMeshStartGuard {
starts: &self.ticket_mesh_starts,
id: id.clone(),
};
let mesh = crate::ticket_mesh::wasm_session::WasmTicketMesh::join(
self.inner.clone(),
&invitation,
action_handler,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.ticket_meshes
.borrow_mut()
.insert(id.clone(), Rc::new(mesh));
Ok(id)
}
pub async fn ticket_mesh_peers(&self, id: String) -> Result<JsValue, JsValue> {
let mesh = self
.ticket_meshes
.borrow()
.get(&id)
.cloned()
.ok_or_else(|| JsValue::from_str("ticket mesh is not active"))?;
serde_wasm_bindgen::to_value(&mesh.peers().await)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub fn ticket_mesh_issuer_node(&self, id: String) -> Result<String, JsValue> {
let mesh = self
.ticket_meshes
.borrow()
.get(&id)
.cloned()
.ok_or_else(|| JsValue::from_str("ticket mesh is not active"))?;
mesh.issuer_node()
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub fn receive_ticket_mesh_frame(
&self,
id: String,
connection_id: String,
frame: String,
) -> Result<(), JsValue> {
let mesh = self
.ticket_meshes
.borrow()
.get(&id)
.cloned()
.ok_or_else(|| JsValue::from_str("ticket mesh is not active"))?;
mesh.receive(connection_id, frame.into_bytes());
Ok(())
}
pub async fn close_ticket_mesh(&self, id: String) -> Result<(), JsValue> {
if self.ticket_mesh_starts.borrow().contains(&id) {
return Err(JsValue::from_str(
"ticket mesh is still starting; await startup before closing",
));
}
let mesh = self.ticket_meshes.borrow().get(&id).cloned();
if let Some(mesh) = mesh {
mesh.close().await;
let mut meshes = self.ticket_meshes.borrow_mut();
if meshes
.get(&id)
.is_some_and(|current| Rc::ptr_eq(current, &mesh))
{
meshes.remove(&id);
}
}
Ok(())
}
pub async fn connect_ticket_invite(&self, invitation: String) -> Result<JsValue, JsValue> {
let result = self
.inner
.connect_ticket_invite(&invitation)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
serde_wasm_bindgen::to_value(&result)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __requireScopePeerAdmission)]
pub fn require_scope_peer_admission(&self, scope: String) -> Result<(), JsValue> {
self.inner
.session_token_registry
.require_scope_peer_admission(&scope)
.map_err(|error| JsValue::from_str(&error))
}
#[wasm_bindgen(js_name = __updateScopePeerAdmission)]
pub fn update_scope_peer_admission(
&self,
scope: String,
revision: f64,
expires_at_ms: f64,
peers_json: String,
) -> Result<(), JsValue> {
let valid_integer = |value: f64| {
value.is_finite()
&& value >= 0.0
&& value <= 9_007_199_254_740_991.0
&& value.fract() == 0.0
};
if !valid_integer(revision) || !valid_integer(expires_at_ms) || peers_json.len() > 4_000
{
return Err(JsValue::from_str("invalid peer admission lease"));
}
let peers = serde_json::from_str(&peers_json)
.map_err(|_| JsValue::from_str("invalid peer admission peers"))?;
self.inner
.session_token_registry
.update_scope_peer_admission(&scope, revision as u64, expires_at_ms as u64, peers)
.map_err(|error| JsValue::from_str(&error))
}
pub fn register_session_token(
&self,
token: String,
grant_scope: String,
max_connections: u32,
) {
self.inner
.register_session_token(token, grant_scope, max_connections);
}
pub fn register_token_until(
&self,
token: String,
grant_scope: String,
max_connections: u32,
expires_at_ms: u64,
) {
self.inner
.register_token_until(token, grant_scope, max_connections, expires_at_ms);
}
#[wasm_bindgen(js_name = setConnectionApplicationCryptoRequired)]
pub fn set_connection_application_crypto_required(
&self,
connection_id: String,
) -> Result<(), JsValue> {
self.inner
.set_connection_application_crypto_required(&connection_id);
Ok(())
}
#[wasm_bindgen(js_name = setConnectionApplicationCryptoKey)]
pub async fn set_connection_application_crypto_key(
&self,
connection_id: String,
key: Vec<u8>,
) -> Result<(), JsValue> {
if key.len() != crate::application_crypto::APPLICATION_KEY_BYTES {
return Err(JsValue::from_str("application crypto key must be 32 bytes"));
}
let mut key_bytes = [0u8; crate::application_crypto::APPLICATION_KEY_BYTES];
key_bytes.copy_from_slice(&key);
self.inner
.set_connection_application_crypto_key(&connection_id, key_bytes);
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
Ok(())
}
#[wasm_bindgen(js_name = clearConnectionApplicationCryptoKey)]
pub fn clear_connection_application_crypto_key(&self, connection_id: String) {
self.inner
.clear_connection_application_crypto_key(&connection_id);
}
#[wasm_bindgen(js_name = ensureConnectionApplicationCrypto)]
pub async fn ensure_connection_application_crypto(
&self,
connection_id: String,
remote_node_id: String,
timeout_ms: u64,
send_frame: js_sys::Function,
) -> Result<(), JsValue> {
let connection_id = connection_id.trim().to_string();
if connection_id.is_empty() {
return Err(JsValue::from_str("connection id is required"));
}
let remote_node_id = remote_node_id.trim().to_string();
if remote_node_id.is_empty() {
return Err(JsValue::from_str("remote node id is required"));
}
let remote_endpoint_id: iroh::EndpointId = remote_node_id.parse().map_err(|error| {
JsValue::from_str(&format!("invalid remote endpoint id: {error}"))
})?;
self.inner
.ensure_connection_manager_record_before_peer_stream(&remote_endpoint_id)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let timeout_ms = timeout_ms.max(1);
let initial = self
.inner
.connection_manager
.get_by_connection_id(&connection_id)
.await
.ok_or_else(|| JsValue::from_str("connection is not registered"))?;
let initial_endpoint = initial
.endpoint_id
.as_deref()
.or(initial.node_id.as_deref())
.map(ToOwned::to_owned)
.ok_or_else(|| JsValue::from_str("connection has no endpoint"))?;
let initial_transport_stable_id = initial.transport_stable_id;
let initial_transport_generation = initial.transport_generation;
let initial_route_generation = initial.route_generation;
self.inner
.set_connection_application_crypto_required(&connection_id);
let started_at_ms = js_sys::Date::now();
let retry_delays_ms = [0_u64, 250, 750];
for (attempt, delay_ms) in retry_delays_ms.into_iter().enumerate() {
if delay_ms > 0 {
let remaining_ms = timeout_ms
.saturating_sub((js_sys::Date::now() - started_at_ms).max(0.0) as u64);
if remaining_ms == 0 {
break;
}
gloo_timers::future::sleep(std::time::Duration::from_millis(
delay_ms.min(remaining_ms),
))
.await;
}
self.inner
.assert_application_crypto_generation(
&connection_id,
&initial_endpoint,
initial_transport_stable_id,
initial_transport_generation,
initial_route_generation,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
if self
.inner
.connection_application_crypto_key(&connection_id)
.is_some()
&& self.inner.connection_application_crypto_is_confirmed(
&connection_id,
initial_transport_stable_id,
)
{
let confirmed_transport_stable_id =
initial_transport_stable_id.ok_or_else(|| {
JsValue::from_str("confirmed connection has no stable ID")
})?;
if !self
.inner
.confirm_managed_connection_readiness_from_transport_proof(
&connection_id,
confirmed_transport_stable_id,
)
.await
{
return Err(JsValue::from_str(
"confirmed application route belongs to a retired generation",
));
}
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
return Ok(());
}
let frame = self
.inner
.application_key_handshake_frame(&connection_id, "capability-update")
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let frame = serde::Serialize::serialize(
&frame,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let send_result = send_frame.call1(&JsValue::NULL, &frame)?;
wasm_bindgen_futures::JsFuture::from(js_sys::Promise::resolve(&send_result))
.await?;
let attempt_deadline_ms = ((attempt + 1) as u64 * timeout_ms / 3).max(1);
loop {
if self
.inner
.connection_application_crypto_key(&connection_id)
.is_some()
&& self.inner.connection_application_crypto_is_confirmed(
&connection_id,
initial_transport_stable_id,
)
{
let confirmed_transport_stable_id = initial_transport_stable_id
.ok_or_else(|| {
JsValue::from_str("confirmed connection has no stable ID")
})?;
if !self
.inner
.confirm_managed_connection_readiness_from_transport_proof(
&connection_id,
confirmed_transport_stable_id,
)
.await
{
return Err(JsValue::from_str(
"confirmed application route belongs to a retired generation",
));
}
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
return Ok(());
}
let elapsed_ms = (js_sys::Date::now() - started_at_ms).max(0.0) as u64;
if elapsed_ms >= attempt_deadline_ms || elapsed_ms >= timeout_ms {
break;
}
self.inner
.assert_application_crypto_generation(
&connection_id,
&initial_endpoint,
initial_transport_stable_id,
initial_transport_generation,
initial_route_generation,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
gloo_timers::future::sleep(std::time::Duration::from_millis(25)).await;
}
}
let current_manager_transport_stable_id = self
.inner
.connection_manager
.get_by_connection_id(&connection_id)
.await
.and_then(|record| record.transport_stable_id);
let current_physical_transport_stable_id = self
.inner
.get_connection(remote_endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let confirmation = self
.inner
.connection_application_crypto_confirmation_snapshot(&connection_id);
return Err(JsValue::from_str(&format!(
"connection {} reciprocal application key agreement timed out after {}ms: frozen_stable_id={:?} current_manager_stable_id={:?} current_physical_stable_id={:?} confirmation={:?}",
connection_id,
timeout_ms,
initial_transport_stable_id,
current_manager_transport_stable_id,
current_physical_transport_stable_id,
confirmation,
)));
}
#[wasm_bindgen(js_name = handleConnectionApplicationCryptoHandshake)]
pub async fn handle_connection_application_crypto_handshake(
&self,
connection_id: String,
remote_node_id: String,
transport_stable_id: Option<u64>,
action: Option<String>,
remote_public_key: Vec<u8>,
) -> Result<JsValue, JsValue> {
if remote_public_key.len() != crate::key_agreement::PUBLIC_KEY_BYTES {
return Err(JsValue::from_str(
"application key agreement public key must be 32 bytes",
));
}
let mut remote_public = [0_u8; crate::key_agreement::PUBLIC_KEY_BYTES];
remote_public.copy_from_slice(&remote_public_key);
let outcome = self
.inner
.accept_application_key_handshake(
&connection_id,
&remote_node_id,
transport_stable_id,
action.as_deref(),
remote_public,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let reply_frame = outcome
.reply_action
.map(|reply_action| {
self.inner
.application_key_handshake_frame(&connection_id, reply_action)
})
.transpose()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
serde::Serialize::serialize(
&reply_frame,
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = connectionApplicationCryptoReady)]
pub fn connection_application_crypto_ready(
&self,
connection_id: String,
transport_stable_id: u64,
) -> bool {
self.inner
.connection_application_crypto_key(&connection_id)
.is_some()
&& self.inner.connection_application_crypto_is_confirmed(
&connection_id,
Some(transport_stable_id),
)
}
#[wasm_bindgen(js_name = protectConnectionApplicationPayload)]
pub fn protect_connection_application_payload(
&self,
connection_id: String,
type_id: u8,
payload: Vec<u8>,
) -> Result<Vec<u8>, JsValue> {
self.inner
.protect_outbound_application_payload_with_type(&connection_id, type_id, &payload)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = openConnectionApplicationPayload)]
pub fn open_connection_application_payload(
&self,
connection_id: String,
type_id: u8,
payload: Vec<u8>,
) -> Result<Vec<u8>, JsValue> {
self.inner
.open_inbound_application_payload_with_type(&connection_id, type_id, &payload)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub fn validate_session_token(&self, token: String) -> Result<String, JsValue> {
self.inner
.validate_session_token(&token)
.map_err(|e| JsValue::from_str(&e))
}
pub async fn validate_connection_token(
&self,
token: String,
connection_id: String,
token_payload: Option<String>,
) -> Result<String, JsValue> {
self.inner
.validate_connection_token(&token, &connection_id, token_payload.as_deref())
.await
.map_err(|e| JsValue::from_str(&e))
}
pub async fn present_session_token(
&self,
endpoint_id: String,
token: String,
token_payload: Option<String>,
device_id: Option<String>,
) -> Result<String, JsValue> {
crate::console_log!(
"[OpenRTC][session-admission][wasm-present] endpoint_id={} claimed_local_device_id={}",
endpoint_id,
device_id.as_deref().unwrap_or("<none>")
);
let endpoint_id_parsed: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
let local_node_id = self.inner.current_node_id().await.ok_or_else(|| {
JsValue::from_str("missing local node id for session-token presentation")
})?;
let connection_id =
crate::client::Client::deterministic_connection_id(&local_node_id, &endpoint_id);
let approval_scope = self
.inner
.present_and_accept_session_token_with_local_claim(
endpoint_id_parsed,
&connection_id,
&token,
token_payload.as_deref(),
None,
device_id,
)
.await
.map_err(|e| JsValue::from_str(&e))?;
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
Ok(approval_scope)
}
pub async fn is_remote_admitted(&self, endpoint_ticket: String) -> Result<bool, JsValue> {
self.inner
.is_remote_admitted(&endpoint_ticket)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[allow(clippy::too_many_arguments)]
pub async fn prepare_reciprocal_admission(
&self,
endpoint_id: String,
expected_transport_stable_id: u64,
stream_instance_id: String,
presentation_id: String,
token: String,
token_payload: String,
device_id: String,
stream_contract: String,
) -> Result<bool, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|error| JsValue::from_str(&format!("{error}")))?;
let stream_contract = match stream_contract.trim() {
"one-shot-admission" => {
crate::native_protocol::TokenStreamContract::OneShotAdmission
}
"persistent-control" => {
crate::native_protocol::TokenStreamContract::PersistentControl
}
other => {
return Err(JsValue::from_str(&format!(
"unsupported reciprocal stream contract: {other}"
)))
}
};
self.inner
.prepare_reciprocal_admission(
endpoint_id,
expected_transport_stable_id,
stream_instance_id.as_str(),
presentation_id.as_str(),
token.as_str(),
token_payload.as_str(),
device_id.as_str(),
stream_contract,
)
.await
.map_err(|error| JsValue::from_str(&error))?;
Ok(true)
}
pub async fn confirm_reciprocal_admission(
&self,
endpoint_id: String,
expected_transport_stable_id: u64,
stream_instance_id: String,
presentation_id: String,
accepted: bool,
approval_scope: Option<String>,
) -> Result<bool, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|error| JsValue::from_str(&format!("{error}")))?;
self.inner
.confirm_reciprocal_admission(
endpoint_id,
expected_transport_stable_id,
stream_instance_id.as_str(),
presentation_id.as_str(),
accepted,
approval_scope.as_deref(),
)
.await
.map_err(|error| JsValue::from_str(&error))?;
self.inner
.emit_current_wasm_connection_state(
&crate::client::Client::deterministic_connection_id(
&self.inner.current_node_id().await.ok_or_else(|| {
JsValue::from_str(
"missing local node id after reciprocal admission ACK",
)
})?,
&endpoint_id.to_string(),
),
)
.await;
Ok(true)
}
pub fn revoke_session_token(&self, token: String) -> Result<JsValue, JsValue> {
serde_wasm_bindgen::to_value(&self.inner.revoke_session_token(&token))
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn revoke_tokens_by_scope(
&self,
grant_scope: String,
) -> Result<JsValue, JsValue> {
let affected = self.inner.begin_revoke_tokens_by_scope(&grant_scope);
for connection_id in &affected {
#[cfg(not(any(feature = "transport-webrtc", feature = "transport-moq")))]
let _ = connection_id;
#[cfg(feature = "transport-webrtc")]
self.retire_iroh_webrtc_carrier(
connection_id.clone(),
Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED.to_string()),
)
.await;
#[cfg(feature = "transport-moq")]
self.retire_iroh_moq_carrier(
connection_id.clone(),
Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED.to_string()),
)
.await;
}
self.inner
.finish_revoke_tokens_by_scope(&grant_scope, &affected)
.await;
serde_wasm_bindgen::to_value(&affected).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub fn clear_session_tokens(&self) {
self.inner.clear_session_tokens();
}
pub fn endpoint_id_from_ticket(&self, ticket: String) -> Result<String, JsValue> {
let (iroh_ticket, _token_suffix) = split_ticket(ticket.trim());
let parsed = EndpointTicket::from_str(iroh_ticket)
.map_err(|e| JsValue::from_str(&format!("Invalid endpoint ticket: {}", e)))?;
Ok(parsed.endpoint_addr().id.to_string())
}
pub async fn disconnect(&self, endpoint_id: String) -> Result<(), JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
let node_guard = self.inner.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.disconnect(endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
} else {
Err(JsValue::from_str("Iroh node not initialized"))
}
}
pub async fn disconnect_transient(&self, endpoint_id: String) -> Result<(), JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
self.inner
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
self.inner.wake_browser_auto_connect();
Ok(())
}
pub async fn is_connected(&self, endpoint_id: String) -> Result<bool, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
Ok(self.inner.is_connected(endpoint_id).await)
}
pub fn runtime_policy(&self) -> Result<JsValue, JsValue> {
serde_wasm_bindgen::to_value(&self.inner.policy_snapshot())
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn add_peer_scope(&self, id: String, scope: String) -> Result<JsValue, JsValue> {
let scopes = self.inner.add_peer_scope(&id, &scope).await;
serde_wasm_bindgen::to_value(&scopes).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn release_peer_scope(
&self,
id: String,
scope: Option<String>,
) -> Result<JsValue, JsValue> {
let scopes = self.inner.release_peer_scope(&id, scope.as_deref()).await;
serde_wasm_bindgen::to_value(&scopes).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn peer_scopes(&self, id: String) -> Result<JsValue, JsValue> {
let scopes = self.inner.peer_scopes(&id).await;
serde_wasm_bindgen::to_value(&scopes).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn same_peer(&self, left: String, right: String) -> Result<bool, JsValue> {
Ok(self.inner.same_peer(&left, &right).await)
}
pub async fn peer_snapshot(&self, id: String) -> Result<JsValue, JsValue> {
let snapshot = self.inner.peer_snapshot(&id).await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn peer_session(&self, id: String) -> Result<JsValue, JsValue> {
let snapshot = self.inner.peer_session(&id).await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn peer_sessions(&self) -> Result<JsValue, JsValue> {
let snapshots = self.inner.peer_sessions().await;
serde_wasm_bindgen::to_value(&snapshots).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn connection_state(&self, connection_id: String) -> Result<JsValue, JsValue> {
let snapshot = self.inner.connection_state(&connection_id).await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
async fn revisit_preferred_iroh_carrier_after_settlement(&self, connection_id: &str) {
let _lifecycle = self.wasm_carrier_peer_lifecycle.lock().await;
let Some(record) = self
.inner
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return;
};
if matches!(
record.active_transport.as_str(),
crate::transport_label::WEBRTC | crate::transport_label::MOQ
) {
return;
}
let Some(capabilities) = self
.wasm_remote_carrier_capabilities
.borrow()
.get(connection_id)
.copied()
else {
return;
};
let Some(remote_endpoint_id) = self
.inner
.wasm_iroh_carrier_remote_endpoint_id(connection_id)
.await
else {
return;
};
if let Err(error) = self
.begin_preferred_iroh_carrier_locked(
connection_id.to_string(),
remote_endpoint_id,
capabilities.webrtc,
capabilities.moq,
)
.await
{
web_sys::console::warn_1(&JsValue::from_str(&format!(
"[OpenRTC][Iroh carrier] settled-generation revisit failed connection_id={} error={:?}",
connection_id, error,
)));
}
}
pub async fn connection_states(&self) -> Result<JsValue, JsValue> {
let snapshots = self.inner.connection_states().await;
serde_wasm_bindgen::to_value(&snapshots).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn wait_for_peer(
&self,
id: String,
timeout_ms: Option<u32>,
) -> Result<JsValue, JsValue> {
let snapshot = self
.inner
.wait_for_peer(&id, timeout_ms.map(|value| value as u64))
.await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn wait_for_settled_scope(
&self,
scope: String,
timeout_ms: Option<u32>,
) -> Result<JsValue, JsValue> {
let snapshot = self
.inner
.wait_for_settled_scope(&scope, timeout_ms.map(u64::from))
.await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn resolve_peer_connection_records(
&self,
id: String,
) -> Result<JsValue, JsValue> {
let records = self.inner.resolve_peer_connection_records(&id).await;
serde_wasm_bindgen::to_value(&records).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn list_managed_connections(&self) -> Result<JsValue, JsValue> {
let records = self.inner.list_managed_connections().await;
serde_wasm_bindgen::to_value(&records).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn bind_connection_device_id(
&self,
connection_id: String,
device_id: String,
) -> Result<JsValue, JsValue> {
let snapshot = self
.inner
.bind_connection_device_id(&connection_id, &device_id)
.await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn acknowledge_session_token_response(
&self,
connection_id: String,
remote_node_id: String,
transport_stable_id: u64,
claimed_device_id: Option<String>,
) -> bool {
self.inner
.acknowledge_wasm_session_token_response(
&connection_id,
&remote_node_id,
transport_stable_id,
claimed_device_id.as_deref(),
)
.await
}
pub async fn bind_node_device_id(
&self,
node_id: String,
device_id: String,
) -> Result<(), JsValue> {
self.inner.bind_node_device_id(&node_id, &device_id).await;
Ok(())
}
pub async fn reject_connection_admission(
&self,
connection_id: String,
reason: String,
) -> Result<JsValue, JsValue> {
self.inner
.reject_session_connection(&connection_id, &reason);
self.inner
.emit_current_wasm_connection_state(&connection_id)
.await;
let snapshot = self.inner.connection_state(&connection_id).await;
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn report_managed_connection_settled(
&self,
connection_id: String,
settled: bool,
device_id: Option<String>,
transport_stable_id: Option<u64>,
transport_generation: Option<u64>,
route_generation: Option<u64>,
) -> Result<JsValue, JsValue> {
let snapshot = match (transport_stable_id, transport_generation, route_generation) {
(Some(transport_stable_id), Some(transport_generation), Some(route_generation)) => {
self.inner
.report_managed_connection_settled_for_transport(
&connection_id,
settled,
transport_stable_id,
transport_generation,
route_generation,
)
.await
}
_ => None,
};
let _ = device_id;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
if settled && snapshot.is_some() {
self.revisit_preferred_iroh_carrier_after_settlement(&connection_id)
.await;
}
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn report_transport_status(
&self,
connection_id: String,
active_transport: String,
parallel_transport: Option<String>,
transport_stable_id: Option<u64>,
transport_generation: Option<u64>,
route_generation: Option<u64>,
) -> Result<JsValue, JsValue> {
let snapshot = match (transport_stable_id, transport_generation, route_generation) {
(Some(transport_stable_id), Some(transport_generation), Some(route_generation)) => {
self.inner
.report_transport_status_for_generation(
&connection_id,
&active_transport,
parallel_transport.as_deref(),
transport_stable_id,
transport_generation,
route_generation,
)
.await
}
_ => None,
};
serde_wasm_bindgen::to_value(&snapshot).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn is_current_transport_stable_id(
&self,
endpoint_id: String,
transport_stable_id: u64,
) -> Result<bool, JsValue> {
let endpoint_id = endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
Ok(self
.inner
.is_current_transport_stable_id(endpoint_id, transport_stable_id)
.await)
}
pub async fn open_bi(&self, endpoint_id: String) -> Result<BiStream, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
self.inner
.assert_raw_peer_stream_allowed(&endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
self.inner
.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
let (send, recv) = self
.inner
.open_bi_internal(endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(BiStream::from_parts(send, recv, endpoint_id.to_string()))
}
pub async fn open_native_main_control_bi(
&self,
endpoint_id: String,
) -> Result<BiStream, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
self.inner
.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
let (transport_stable_id, send, recv) = self
.inner
.open_bi_internal_with_transport_stable_id(endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(BiStream::from_parts_with_transport_stable_id(
send,
recv,
endpoint_id.to_string(),
transport_stable_id,
))
}
pub async fn send_native_main_control_frame(
&self,
endpoint_id: String,
frame: Vec<u8>,
) -> Result<(), JsValue> {
send_native_signal_control_frame(&self.inner, endpoint_id, frame).await
}
pub async fn open_peer_bi(
&self,
id: String,
timeout_ms: Option<u32>,
) -> Result<BiStream, JsValue> {
let (_connection_id, remote_node_id, send, recv) = self
.inner
.open_peer_bi(&id, timeout_ms.map(|value| value as u64))
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(BiStream::from_peer_parts(send, recv, remote_node_id))
}
pub async fn send_peer_application_frame(
&self,
id: String,
frame: Vec<u8>,
timeout_ms: Option<u32>,
) -> Result<(), JsValue> {
self.inner
.send_peer_application_frame(&id, &frame, timeout_ms.map(|value| value as u64))
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn open_peer_bi_explicit_file_sender(
&self,
id: String,
timeout_ms: Option<u32>,
) -> Result<PeerUniStream, JsValue> {
let (_connection_id, _remote_node_id, send) = self
.inner
.open_peer_bi_explicit_file_sender(&id, timeout_ms.map(|value| value as u64))
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(peer_uni_stream_from_send(send))
}
pub async fn open_peer_bi_transport_only(
&self,
id: String,
timeout_ms: Option<u32>,
) -> Result<BiStream, JsValue> {
let (_connection_id, remote_node_id, send, recv) = self
.inner
.open_peer_bi_transport_only(&id, timeout_ms.map(|value| value as u64))
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(BiStream::from_parts(send, recv, remote_node_id))
}
pub async fn open_uni(&self, endpoint_id: String) -> Result<PeerUniStream, JsValue> {
let endpoint_id: iroh::EndpointId = endpoint_id
.parse()
.map_err(|e| JsValue::from_str(&format!("{}", e)))?;
self.inner
.assert_raw_peer_stream_allowed(&endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
self.inner
.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
let node_guard = self.inner.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
let send = node
.open_uni(endpoint_id.clone())
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(peer_uni_stream_from_send(
crate::application_crypto_streams::PeerSendStream::plain(send),
))
} else {
Err(JsValue::from_str("Iroh node not initialized"))
}
}
pub async fn open_peer_uni(
&self,
id: String,
timeout_ms: Option<u32>,
) -> Result<PeerUniStream, JsValue> {
let (_connection_id, _remote_node_id, send) = self
.inner
.open_peer_uni(&id, timeout_ms.map(|value| value as u64))
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
Ok(peer_uni_stream_from_send(send))
}
#[wasm_bindgen(js_name = sendPeerDatagram)]
pub async fn send_peer_datagram(
&self,
id: String,
payload: Vec<u8>,
max_age_ms: Option<u32>,
) -> Result<bool, JsValue> {
match max_age_ms {
Some(max_age_ms) => {
self.inner
.send_peer_datagram_with_max_age(&id, &payload, u64::from(max_age_ms))
.await
}
None => self
.inner
.send_peer_datagram(&id, &payload)
.await
.map(|()| true),
}
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = receivePeerDatagram)]
pub async fn receive_peer_datagram(&self, id: String) -> Result<Vec<u8>, JsValue> {
self.inner
.receive_peer_datagram(&id)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = peerDatagramMaxSize)]
pub fn peer_datagram_max_size(&self) -> usize {
Client::MAX_PEER_DATAGRAM_BYTES
}
#[wasm_bindgen(js_name = createPeerDatagramPolicy)]
pub fn create_peer_datagram_policy(&self) -> WasmPeerDatagramPolicy {
WasmPeerDatagramPolicy {
inner: RefCell::new(crate::datagrams::PeerDatagramPolicy::new(
Client::MAX_PEER_DATAGRAM_BYTES,
)),
}
}
pub async fn send_peer(&self, id: String, data: Vec<u8>) -> Result<(), JsValue> {
self.inner
.send_peer(&id, &data)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn iroh_path_kind(&self, peer_id: String) -> String {
match self.inner.iroh_path_kind(&peer_id).await {
crate::client::IrohPathKind::DirectQuic => "direct-quic".to_string(),
crate::client::IrohPathKind::DirectLan => "direct-lan".to_string(),
crate::client::IrohPathKind::Relay => "relay".to_string(),
crate::client::IrohPathKind::Ble => "ble".to_string(),
crate::client::IrohPathKind::WebRtc => "webrtc".to_string(),
crate::client::IrohPathKind::Moq => "moq".to_string(),
crate::client::IrohPathKind::Unknown => "unknown".to_string(),
}
}
pub async fn iroh_transport_rtt_ms(&self, peer_id: String) -> Option<u32> {
self.inner
.iroh_transport_rtt_ms(&peer_id)
.await
.map(|value| value.min(u32::MAX as u64) as u32)
}
pub async fn incoming_streams(&self) -> Result<JsReadableStream, JsValue> {
let (stream_owner, stream) = {
let node_guard = self.inner.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
(self.inner.clone(), node.incoming_streams_stream())
} else {
return Err(JsValue::from_str("Iroh node not initialized"));
}
};
use futures::StreamExt;
let mapped_stream = stream.filter_map(move |incoming| {
let stream_owner = stream_owner.clone();
async move {
let owned = stream_owner
.incoming_stream_generation_is_owned(
incoming.endpoint_id,
incoming.transport_stable_id,
)
.await;
let admission_candidate = !owned
&& matches!(
&incoming.stream,
crate::wasm_node::IncomingStreamType::Bi(_, _)
)
&& stream_owner
.incoming_stream_generation_is_authenticated_reconnect_candidate(
incoming.endpoint_id,
incoming.transport_stable_id,
)
.await;
(owned || admission_candidate).then(|| {
crate::wasm_node::BiStream::incoming_to_js_value(
incoming,
admission_candidate,
)
})
}
});
Ok(wasm_streams::ReadableStream::from_stream(mapped_stream).into_raw())
}
#[wasm_bindgen(js_name = inspectIncomingStreamPrefix)]
pub fn inspect_incoming_stream_prefix(
&self,
bytes: Vec<u8>,
reached_eof: bool,
) -> Result<JsValue, JsValue> {
serde::Serialize::serialize(
&crate::stream_metadata::inspect_prefix_for_binding(&bytes, reached_eof),
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = decodeIncomingChannelEnvelope)]
pub fn decode_incoming_channel_envelope(&self, bytes: Vec<u8>) -> Result<JsValue, JsValue> {
serde::Serialize::serialize(
&crate::stream_metadata::decode_prefix_for_binding(&bytes),
&serde_wasm_bindgen::Serializer::json_compatible(),
)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn update_presence(
&self,
user_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
ttl_ms: Option<u64>,
) -> Result<(), JsValue> {
self.inner
.update_presence_with_ttl(
&user_id,
&device_name,
&ticket,
ttl_ms.unwrap_or(300_000),
metadata.as_deref(),
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn send_message(
&self,
target_id: String,
payload: String,
state: Option<String>,
reply_payload: Option<String>,
) -> Result<String, JsValue> {
self.inner
.send_message(
&target_id,
&payload,
state.as_deref(),
reply_payload.as_deref(),
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn set_offline(&self, user_id: String) -> Result<(), JsValue> {
self.inner
.set_offline(&user_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn update_device(
&self,
user_id: String,
device_id: String,
device_name: Option<String>,
capabilities: Option<JsValue>,
metadata: Option<String>,
) -> Result<(), JsValue> {
let parsed_capabilities = match capabilities {
Some(value) if !value.is_null() && !value.is_undefined() => Some(
serde_wasm_bindgen::from_value::<crate::signaling::DeviceCapabilities>(value)
.map_err(|e| JsValue::from_str(&e.to_string()))?,
),
_ => None,
};
self.inner
.update_device(
&user_id,
&device_id,
device_name.as_deref(),
parsed_capabilities,
metadata.as_deref(),
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn delete_device(
&self,
user_id: String,
device_id: String,
) -> Result<(), JsValue> {
self.inner
.delete_device(&user_id, &device_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub fn force_reconnect_snapshot(&self) {
self.inner.clone().force_reconnect_snapshot();
}
pub fn stop_presence_loop(&self) {
self.inner.stop_presence_loop();
}
pub fn stop_auto_connect(&self) {
self.inner.stop_auto_connect();
self.inner.stop_browser_auto_connect();
}
pub fn start_auto_connect(
&self,
user_id: String,
local_device_id: String,
) -> Result<(), JsValue> {
self.inner
.start_browser_auto_connect(user_id, local_device_id)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[deprecated(note = "use submit_browser_desired_peers_with_sparse_fanout for sparse rooms")]
pub fn submit_browser_desired_peers(
&self,
revision: u32,
peers_json: String,
) -> Result<bool, JsValue> {
self.inner
.submit_browser_desired_peers(u64::from(revision), &peers_json)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub async fn submit_browser_desired_peers_with_sparse_fanout(
&self,
revision: u32,
peers_json: String,
capability: Option<String>,
local_device_id: Option<String>,
local_device_key_x: Option<String>,
sparse: Option<bool>,
) -> Result<bool, JsValue> {
let peers_json = if sparse.unwrap_or(false) {
let capability = capability
.as_deref()
.ok_or_else(|| JsValue::from_str("sparse fanout requires a capability"))?;
let local_device_id = local_device_id
.as_deref()
.ok_or_else(|| JsValue::from_str("sparse fanout requires a local device id"))?;
let local_device_key_x = local_device_key_x.as_deref().ok_or_else(|| {
JsValue::from_str("sparse fanout requires a local device key")
})?;
self.inner
.configure_sparse_fanout(
capability,
u64::from(revision),
local_device_id,
local_device_key_x,
&peers_json,
true,
)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?
.0
} else {
peers_json
};
self.inner
.submit_browser_desired_peers(u64::from(revision), &peers_json)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = prepareSparseFanoutMessage)]
pub async fn prepare_sparse_fanout_message(
&self,
capability: String,
payload: Vec<u8>,
) -> Result<JsValue, JsValue> {
let request = self
.inner
.prepare_sparse_fanout_message(&capability, &payload)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
serde_wasm_bindgen::to_value(&request)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = finalizeSparseFanoutMessage)]
pub async fn finalize_sparse_fanout_message(
&self,
request_id: String,
signature: String,
) -> Result<Vec<u8>, JsValue> {
self.inner
.finalize_sparse_fanout_message(&request_id, &signature)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = acceptSparseFanoutMessage)]
pub async fn accept_sparse_fanout_message(
&self,
capability: String,
source_peer_id: String,
encoded: Vec<u8>,
) -> Result<JsValue, JsValue> {
let decision = self
.inner
.accept_sparse_fanout_message(&capability, &source_peer_id, &encoded)
.await
.map_err(|error| JsValue::from_str(&error.to_string()))?;
serde_wasm_bindgen::to_value(&decision)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = sparseFanoutDiagnostics)]
pub async fn sparse_fanout_diagnostics(
&self,
capability: String,
) -> Result<JsValue, JsValue> {
serde_wasm_bindgen::to_value(&self.inner.sparse_fanout_diagnostics(&capability).await)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = recordSparseFanoutForwardQueueDrop)]
pub async fn record_sparse_fanout_forward_queue_drop(
&self,
capability: String,
count: u32,
) -> Result<(), JsValue> {
self.inner
.record_sparse_fanout_forward_queue_drop(&capability, u64::from(count))
.await
.map_err(|error| JsValue::from_str(&error.to_string()))
}
pub fn wake_browser_auto_connect(&self) -> bool {
self.inner.wake_browser_auto_connect()
}
pub async fn set_auto_connect_excluded(&self, device_id: String, excluded: bool) {
if excluded {
self.inner.exclude_peer_and_publish(&device_id).await;
} else {
self.inner.unexclude_peer_and_publish(&device_id).await;
}
self.inner.wake_browser_auto_connect();
}
pub fn is_auto_connect_excluded(&self, device_id: String) -> bool {
self.inner.is_auto_connect_excluded(&device_id)
}
pub async fn disconnect_device(
&self,
device_id: String,
node_id_hint: Option<String>,
) -> Result<JsValue, JsValue> {
let retired = self
.inner
.disconnect_device(&device_id, node_id_hint.as_deref())
.await;
serde_wasm_bindgen::to_value(&retired).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub fn stop_auth_scoped_activity(&self) {
self.inner.stop_auth_scoped_activity();
self.inner.stop_browser_auto_connect();
}
pub fn start_presence_loop(
&self,
user_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
) {
self.inner
.clone()
.start_signaling_loop(user_id, device_name, ticket, metadata);
}
pub async fn search_devices(&self, user_id: String) -> Result<JsValue, JsValue> {
let devices = self
.inner
.search_devices(&user_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_wasm_bindgen::to_value(&devices).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn devices_with_status(&self, user_id: String) -> Result<JsValue, JsValue> {
let devices = self
.inner
.devices_with_status(&user_id)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_wasm_bindgen::to_value(&devices).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn connect_device(
&self,
device_id: Option<String>,
endpoint_ticket: String,
timeout_ms: Option<u32>,
) -> Result<JsValue, JsValue> {
let dial = self
.inner
.connect_device(device_id.as_deref(), &endpoint_ticket);
let result = if let Some(timeout_ms) = timeout_ms {
let deadline = gloo_timers::future::sleep(std::time::Duration::from_millis(
u64::from(timeout_ms.clamp(1, 120_000)),
));
futures::pin_mut!(dial, deadline);
match futures::future::select(dial, deadline).await {
futures::future::Either::Left((result, _)) => result,
futures::future::Either::Right(_) => {
return Err(JsValue::from_str("endpoint ticket dial timed out"));
}
}
} else {
dial.await
}
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_wasm_bindgen::to_value(&result).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn connect_known_device_with_token(
&self,
device_id: String,
token: String,
scope: String,
max_connections: u32,
expires_at_ms: Option<u64>,
lookup_timeout_ms: Option<u32>,
) -> Result<JsValue, JsValue> {
let result = self
.inner
.connect_known_device_with_token(
&device_id,
&token,
&scope,
max_connections,
expires_at_ms,
lookup_timeout_ms.map(u64::from),
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_wasm_bindgen::to_value(&result).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn observe_known_device_endpoint(
&self,
device_id: String,
endpoint_ticket: String,
) -> Result<(), JsValue> {
self.inner
.observe_known_device_endpoint(&device_id, &endpoint_ticket)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn create_session(&self, session_json: String) -> Result<(), JsValue> {
let session: crate::signaling::SignalingSession =
serde_json::from_str(&session_json)
.map_err(|e| JsValue::from_str(&e.to_string()))?;
self.inner
.create_session(session)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn update_session(
&self,
session_id: String,
update_json: String,
) -> Result<(), JsValue> {
let update_data: serde_json::Value = serde_json::from_str(&update_json)
.map_err(|e| JsValue::from_str(&e.to_string()))?;
self.inner
.update_session(&session_id, update_data)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn create_room(
&self,
room_id: String,
user_id: String,
ticket_str: String,
my_node_id: String,
tag: String,
max_members: Option<u32>,
) -> Result<bool, JsValue> {
self.inner
.room
.create_room(
&room_id,
&user_id,
&ticket_str,
&my_node_id,
&tag,
max_members,
)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn join_room(
&self,
room_id: String,
user_id: String,
ticket_str: String,
my_node_id: String,
tag: String,
) -> Result<(), JsValue> {
self.inner
.room
.join_room(&room_id, &user_id, &ticket_str, &my_node_id, &tag)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn get_members(
&self,
room_id: String,
my_node_id: String,
tag: String,
) -> Result<String, JsValue> {
let members = self
.inner
.room
.get_members(&room_id, &my_node_id, &tag)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))?;
serde_json::to_string(&members).map_err(|e| JsValue::from_str(&e.to_string()))
}
pub async fn leave_room(
&self,
room_id: String,
my_node_id: String,
tag: String,
) -> Result<(), JsValue> {
self.inner
.room
.leave_room(&room_id, &my_node_id, &tag)
.await
.map_err(|e| JsValue::from_str(&e.to_string()))
}
}
#[cfg(feature = "iroh-protocols-wasm")]
#[wasm_bindgen]
impl WasmClient {
#[wasm_bindgen(js_name = __initPersistentIrohProtocols)]
pub async fn init_persistent_iroh_protocols(
&self,
replica_store: JsValue,
) -> Result<(), JsValue> {
let mut guard = self.persistent_protocols.lock().await;
if guard.is_some() {
return Ok(());
}
let node = self
.inner
.iroh_node
.read()
.await
.as_ref()
.cloned()
.ok_or_else(|| {
JsValue::from_str("OpenRTC endpoint must be initialized before protocols")
})?;
let store = crate::wasm_docs_persistence::JsReplicaStore::new(replica_store)
.map_err(|error| JsValue::from_str(&error.to_string()))?;
let actor = crate::wasm_docs_persistence::WasmPersistentDocsActor::hydrate(
store,
node.endpoint().clone(),
)
.await
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
node.install_standard_protocols(
actor.docs_protocol(),
actor.blobs_protocol(),
actor.gossip_protocol(),
)
.await
.map_err(|error| JsValue::from_str(&format!("{error:#}")))?;
*guard = Some(actor);
Ok(())
}
#[wasm_bindgen(js_name = __reconcilePersistentIrohProtocols)]
pub async fn reconcile_persistent_iroh_protocols(&self) -> Result<u32, JsValue> {
let guard = self.persistent_protocols.lock().await;
let Some(actor) = guard.as_ref() else {
return Ok(0);
};
actor.reconcile_sync().await.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __importPersistentIrohAuthor)]
pub async fn import_persistent_iroh_author(
&self,
author_secret: Vec<u8>,
) -> Result<String, JsValue> {
let guard = self.persistent_protocols.lock().await;
let actor = persistent_protocols(&guard)?;
actor
.import_author(author_secret)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __importPersistentIrohNamespace)]
pub async fn import_persistent_iroh_namespace(
&self,
capability_kind: String,
capability: Vec<u8>,
generation: u64,
share_revision: u64,
) -> Result<String, JsValue> {
let mut guard = self.persistent_protocols.lock().await;
let actor = persistent_protocols_mut(&mut guard)?;
actor
.import_namespace(
capability_kind.as_str(),
capability,
generation,
share_revision,
)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __createPersistentIrohNamespace)]
pub async fn create_persistent_iroh_namespace(
&self,
generation: u64,
share_revision: u64,
) -> Result<JsValue, JsValue> {
let mut guard = self.persistent_protocols.lock().await;
let descriptor = persistent_protocols_mut(&mut guard)?
.create_namespace(generation, share_revision)
.await
.map_err(js_protocol_error)?;
serde_wasm_bindgen::to_value(&descriptor)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __importPersistentIrohTicket)]
pub async fn import_persistent_iroh_ticket(
&self,
ticket: String,
generation: u64,
share_revision: u64,
) -> Result<String, JsValue> {
let mut guard = self.persistent_protocols.lock().await;
persistent_protocols_mut(&mut guard)?
.import_ticket(&ticket, generation, share_revision)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __sharePersistentIrohNamespace)]
pub async fn share_persistent_iroh_namespace(
&self,
namespace_id: String,
writable: bool,
) -> Result<String, JsValue> {
let guard = self.persistent_protocols.lock().await;
persistent_protocols(&guard)?
.share(&namespace_id, writable)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __removePersistentIrohNamespace)]
pub async fn remove_persistent_iroh_namespace(
&self,
namespace_id: String,
generation: u64,
share_revision: u64,
) -> Result<(), JsValue> {
let mut guard = self.persistent_protocols.lock().await;
persistent_protocols_mut(&mut guard)?
.remove_namespace(&namespace_id, generation, share_revision)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __putPersistentIrohBytes)]
pub async fn put_persistent_iroh_bytes(
&self,
namespace_id: String,
key: Vec<u8>,
value: Vec<u8>,
) -> Result<JsValue, JsValue> {
let guard = self.persistent_protocols.lock().await;
let receipt = persistent_protocols(&guard)?
.set_bytes(&namespace_id, key, value)
.await
.map_err(js_protocol_error)?;
serde_wasm_bindgen::to_value(&receipt)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __setPersistentIrohHash)]
pub async fn set_persistent_iroh_hash(
&self,
namespace_id: String,
key: Vec<u8>,
content_hash: String,
content_length: u64,
) -> Result<String, JsValue> {
let guard = self.persistent_protocols.lock().await;
persistent_protocols(&guard)?
.set_hash(&namespace_id, key, &content_hash, content_length)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __deletePersistentIrohPrefix)]
pub async fn delete_persistent_iroh_prefix(
&self,
namespace_id: String,
prefix: Vec<u8>,
) -> Result<JsValue, JsValue> {
let guard = self.persistent_protocols.lock().await;
let receipt = persistent_protocols(&guard)?
.delete_prefix(&namespace_id, prefix)
.await
.map_err(js_protocol_error)?;
serde_wasm_bindgen::to_value(&receipt)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __queryPersistentIrohNamespace)]
pub async fn query_persistent_iroh_namespace(
&self,
namespace_id: String,
key_prefix: Vec<u8>,
) -> Result<JsValue, JsValue> {
let guard = self.persistent_protocols.lock().await;
let entries = persistent_protocols(&guard)?
.query(&namespace_id, key_prefix)
.await
.map_err(js_protocol_error)?;
serde_wasm_bindgen::to_value(&entries)
.map_err(|error| JsValue::from_str(&error.to_string()))
}
#[wasm_bindgen(js_name = __hydratePersistentIrohBlob)]
pub async fn hydrate_persistent_iroh_blob(
&self,
content_hash: String,
) -> Result<(), JsValue> {
let guard = self.persistent_protocols.lock().await;
persistent_protocols(&guard)?
.hydrate_blob(&content_hash)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __acknowledgePersistentIrohOutbox)]
pub async fn acknowledge_persistent_iroh_outbox(
&self,
operation_id: String,
) -> Result<(), JsValue> {
let guard = self.persistent_protocols.lock().await;
persistent_protocols(&guard)?
.acknowledge_outbox(&operation_id)
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __flushPersistentIrohProtocols)]
pub async fn flush_persistent_iroh_protocols(&self) -> Result<(), JsValue> {
let guard = self.persistent_protocols.lock().await;
persistent_protocols(&guard)?
.flush()
.await
.map_err(js_protocol_error)
}
#[wasm_bindgen(js_name = __shutdownPersistentIrohProtocols)]
pub async fn shutdown_persistent_iroh_protocols(&self) -> Result<(), JsValue> {
if let Some(actor) = self.persistent_protocols.lock().await.take() {
actor.shutdown().await.map_err(js_protocol_error)?;
}
Ok(())
}
}
#[cfg(feature = "iroh-protocols-wasm")]
fn persistent_protocols(
guard: &Option<crate::wasm_docs_persistence::WasmPersistentDocsActor>,
) -> Result<&crate::wasm_docs_persistence::WasmPersistentDocsActor, JsValue> {
guard
.as_ref()
.ok_or_else(|| JsValue::from_str("persistent Iroh protocols are not initialized"))
}
#[cfg(feature = "iroh-protocols-wasm")]
fn persistent_protocols_mut(
guard: &mut Option<crate::wasm_docs_persistence::WasmPersistentDocsActor>,
) -> Result<&mut crate::wasm_docs_persistence::WasmPersistentDocsActor, JsValue> {
guard
.as_mut()
.ok_or_else(|| JsValue::from_str("persistent Iroh protocols are not initialized"))
}
#[cfg(feature = "iroh-protocols-wasm")]
fn js_protocol_error(error: impl std::fmt::Display) -> JsValue {
JsValue::from_str(&format!("{error:#}"))
}
#[wasm_bindgen(start)]
pub fn start() {
console_error_panic_hook::set_once();
}
}