use super::firebase::client::FirestoreClient;
use super::firebase::room::{FirebaseRoomBackend, RoomBackend};
#[cfg(any(target_os = "ios", target_os = "android"))]
use super::firebase::signaling::FirebaseSignalingBackend;
#[cfg(target_arch = "wasm32")]
use super::firebase::signaling::FirebaseSignalingBackend as SignalingBackendImpl;
#[cfg(all(
not(target_arch = "wasm32"),
not(any(target_os = "ios", target_os = "android"))
))]
use super::firebase::native_signaling::NativeFirestoreSignalingBackend as SignalingBackendImpl;
use crate::signaling::SignalingBackend;
use iroh::Endpoint;
use iroh_tickets::endpoint::EndpointTicket;
use std::collections::{HashMap, HashSet};
use std::str::FromStr;
#[cfg(not(target_arch = "wasm32"))]
use std::sync::atomic::Ordering;
use std::sync::atomic::{AtomicBool, AtomicU64};
use std::sync::{Arc, Mutex, RwLock as StdRwLock};
use tokio::sync::RwLock;
#[cfg(not(target_arch = "wasm32"))]
use crate::client::scope_classifier::{default_scope_classifier, ScopeClassifier};
#[cfg(target_arch = "wasm32")]
fn wasm_init_log(stage: &str) {
let now = js_sys::Date::now();
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[OPENRTC][WASM-INIT] ts_ms={:.0} stage={}",
now, stage
)));
}
#[cfg(target_arch = "wasm32")]
use crate::wasm_node::{AcceptEvent, ConnectEvent, IrohWasmNode};
#[cfg(not(target_arch = "wasm32"))]
use crate::native_node::{AcceptEvent, ConnectEvent, IncomingStream, IrohNativeNode};
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct EndpointHandle {
pub node_id: String,
pub node_addr: String,
}
#[cfg(not(target_arch = "wasm32"))]
pub struct BiStream {
pub send: crate::application_crypto_streams::PeerSendStream,
pub recv: crate::application_crypto_streams::PeerRecvStream,
pub id: String,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone)]
pub struct NativePeerDataEvent {
pub connection_id: String,
pub remote_node_id: Option<String>,
pub transport: String,
pub payload: Vec<u8>,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone)]
pub(crate) struct CachedManagedScopeTicket {
pub scope: crate::session_token::GrantScope,
pub token: String,
pub max_connections: u32,
pub compound_ticket: String,
pub iroh_ticket: String,
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
#[derive(Debug, Clone)]
pub(crate) struct NativeWebRTCSuppression {
pub until_ms: i64,
pub reason: String,
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
#[derive(Debug, Clone)]
pub(crate) struct NativeWebRTCAttemptInFlight {
pub negotiation_id: String,
pub started_at_ms: i64,
pub expires_at_ms: i64,
pub preserve_on_direct_path: bool,
pub last_skip_log_ms: i64,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AppLimits {
pub devices_per_user: i64,
pub max_rooms: i64,
pub max_members_per_room: i64,
}
impl Default for AppLimits {
fn default() -> Self {
Self {
devices_per_user: 2,
max_rooms: 10,
max_members_per_room: 5,
}
}
}
#[allow(dead_code)]
#[derive(Clone)]
pub struct Client {
app_tag: String,
signaling: Arc<dyn SignalingBackend>,
pub room: Arc<dyn RoomBackend>,
node_id: Arc<RwLock<Option<String>>>,
pub iroh_endpoint: Arc<RwLock<Option<Endpoint>>>,
#[cfg(target_arch = "wasm32")]
pub iroh_node: Arc<RwLock<Option<IrohWasmNode>>>,
#[cfg(not(target_arch = "wasm32"))]
pub iroh_node: Arc<RwLock<Option<IrohNativeNode>>>,
pub connection_manager: Arc<crate::connection_manager::ConnectionManager>,
#[cfg(not(target_arch = "wasm32"))]
native_device_identity: Arc<RwLock<Option<crate::native_device::NativeDeviceIdentity>>>,
#[cfg(not(target_arch = "wasm32"))]
native_device_base_dir: Arc<RwLock<Option<std::path::PathBuf>>>,
#[cfg(not(target_arch = "wasm32"))]
native_device_updates:
tokio::sync::broadcast::Sender<crate::native_device::NativeDeviceIdentity>,
#[cfg(not(target_arch = "wasm32"))]
native_connection_state_updates: tokio::sync::broadcast::Sender<ConnectionStateSnapshot>,
#[cfg(not(target_arch = "wasm32"))]
native_peer_data_updates: tokio::sync::broadcast::Sender<NativePeerDataEvent>,
auto_connect_loop_key: Arc<Mutex<Option<(String, String)>>>,
auto_connect_generation: Arc<AtomicU64>,
auto_connect_excluded: Arc<Mutex<HashSet<String>>>,
auto_connect_excluded_node_aliases: Arc<Mutex<HashMap<String, HashSet<String>>>>,
pub session_token_registry: Arc<crate::session_token::SessionTokenRegistry>,
connection_application_crypto_keys:
Arc<StdRwLock<HashMap<String, [u8; crate::application_crypto::APPLICATION_KEY_BYTES]>>>,
connection_application_crypto_required: Arc<StdRwLock<HashSet<String>>>,
connection_application_crypto_outbound_sequences: Arc<StdRwLock<HashMap<String, u64>>>,
connection_application_key_agreements:
Arc<StdRwLock<HashMap<String, crate::key_agreement::EphemeralKeyAgreement>>>,
#[cfg(not(target_arch = "wasm32"))]
managed_scope_tickets: Arc<StdRwLock<HashMap<String, CachedManagedScopeTicket>>>,
app_backgrounded: Arc<AtomicBool>,
#[cfg(target_arch = "wasm32")]
wasm_accept_bridge_started: Arc<std::sync::atomic::AtomicBool>,
#[cfg(target_arch = "wasm32")]
last_emitted_connection_states:
Arc<Mutex<std::collections::HashMap<String, WasmConnectionStateFingerprint>>>,
#[cfg(target_arch = "wasm32")]
last_empty_peer_sessions_warning_ms: Arc<AtomicU64>,
iroh_init_guard: Arc<tokio::sync::Mutex<()>>,
managed_connect_gates: Arc<tokio::sync::Mutex<HashMap<String, Arc<tokio::sync::Mutex<()>>>>>,
pub(crate) presence_loop_tx:
Arc<Mutex<Option<tokio::sync::mpsc::Sender<crate::presence::PresenceCommand>>>>,
#[cfg(not(target_arch = "wasm32"))]
pub app_limits: Arc<RwLock<AppLimits>>,
#[cfg(not(target_arch = "wasm32"))]
pub(crate) project_id: String,
#[cfg(not(target_arch = "wasm32"))]
pub(crate) token_provider: Arc<dyn Fn() -> Option<String> + Send + Sync>,
transport_config: Arc<RwLock<TransportConfig>>,
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "transport-lan", openrtc_unpublished_ble)
))]
local_discovery_registry: crate::local_discovery::LocalDiscoveryRegistry,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
mdns_address_lookup: Arc<RwLock<Option<iroh_mdns_address_lookup::MdnsAddressLookup>>>,
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
ble_transport: Arc<RwLock<Option<Arc<iroh_ble_transport::BleTransport>>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_sessions:
Arc<RwLock<HashMap<String, Arc<crate::transport::NativeWebRTCDataChannel>>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_suppressions: Arc<RwLock<HashMap<String, NativeWebRTCSuppression>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_attempt_counts: Arc<RwLock<HashMap<String, u32>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_start_gates: Arc<tokio::sync::Mutex<HashMap<String, Arc<tokio::sync::Notify>>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
pub(crate) native_webrtc_attempts_in_flight:
Arc<RwLock<HashMap<String, NativeWebRTCAttemptInFlight>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_retry_deadlines: Arc<tokio::sync::Mutex<HashMap<String, i64>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
deferred_managed_retirements: Arc<RwLock<HashMap<String, Option<String>>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
native_moq_sessions: Arc<RwLock<HashMap<String, Arc<crate::transport::NativeMoQSession>>>>,
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
pub(crate) native_signal_streams: Arc<
tokio::sync::Mutex<HashMap<String, Arc<tokio::sync::Mutex<iroh::endpoint::SendStream>>>>,
>,
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
scoped_connection_actor_registry: Arc<
RwLock<Option<Arc<crate::client::scoped_connection_actor::ScopedConnectionActorRegistry>>>,
>,
#[cfg(not(target_arch = "wasm32"))]
auth_readiness: Arc<crate::client::auth_readiness::AuthReadinessStore>,
#[cfg(not(target_arch = "wasm32"))]
scope_classifier: Arc<RwLock<Arc<dyn ScopeClassifier>>>,
}
fn now_millis_i64() -> i64 {
crate::firebase::now_millis_u64().min(i64::MAX as u64) as i64
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IrohPathKind {
DirectQuic,
DirectLan,
Relay,
Ble,
Unknown,
}
impl IrohPathKind {
pub fn is_relay_path(self) -> bool {
matches!(self, Self::Relay)
}
pub fn transport_label(self) -> &'static str {
match self {
Self::DirectQuic => crate::transport_label::IROH_QUIC,
Self::DirectLan => crate::transport_label::IROH_LAN,
Self::Relay => crate::transport_label::IROH_RELAY,
Self::Ble => crate::transport_label::BLE,
Self::Unknown => crate::transport_label::IROH,
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct IrohLanConfig {
#[serde(default = "default_lan_enabled")]
pub enabled: bool,
#[serde(default = "default_lan_advertise")]
pub advertise: bool,
}
fn default_lan_enabled() -> bool {
true
}
fn default_lan_advertise() -> bool {
true
}
impl Default for IrohLanConfig {
fn default() -> Self {
Self {
enabled: true,
advertise: true,
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct BleConfig {
#[serde(default)]
pub enabled: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub connect_timeout_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retry_attempts: Option<u8>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retry_backoff_ms: Option<u64>,
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct BleHardwareHarnessDebugSnapshot {
pub enabled: bool,
pub local_node_id: Option<String>,
pub target_node_id: Option<String>,
pub local_service_uuid: Option<String>,
pub target_service_uuid: Option<String>,
pub target_seen: bool,
pub route_pipes: usize,
pub route_pipe_tombstones: usize,
pub route_scan_hints: usize,
pub route_pending: usize,
pub route_routable: usize,
pub route_reservations: usize,
pub peers: Vec<BleHardwareHarnessPeerSnapshot>,
pub note: Option<String>,
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct BleHardwareHarnessPeerSnapshot {
pub device_id: String,
pub phase: String,
pub phase_detail: Option<String>,
pub consecutive_failures: u32,
pub connect_path: Option<String>,
pub verified_endpoint: Option<String>,
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct BleHardwareHarnessNativeProbe {
pub endpoint_id: String,
pub service_uuid: String,
pub device_id: Option<String>,
pub stages: Vec<BleHardwareHarnessNativeProbeStage>,
pub services: Vec<BleHardwareHarnessNativeProbeService>,
pub success: bool,
pub error: Option<String>,
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct BleHardwareHarnessNativeProbeStage {
pub name: String,
pub ok: bool,
pub elapsed_ms: u128,
pub detail: Option<String>,
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct BleHardwareHarnessNativeProbeService {
pub uuid: String,
pub characteristics: Vec<String>,
}
pub(crate) const MANAGED_SETTLE_DEADLINE_MS: i64 =
crate::runtime_policy::MANAGED_SETTLE_DEADLINE_MS;
#[cfg(not(target_arch = "wasm32"))]
pub(crate) const SESSION_ADMISSION_TIMEOUT_MS: u64 =
crate::runtime_policy::SESSION_ADMISSION_TIMEOUT_MS;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum DeviceConnectionStatus {
Disconnected,
Connecting,
Connected,
Failed,
Closed,
Online,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum DevicePresenceStatus {
Online,
Idle,
Offline,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum ReadinessState {
Connecting,
TransportOnly,
Settling,
Routable,
AwaitingReplacement,
Closed,
Failed,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct DeviceStatusSnapshot {
#[serde(flatten)]
pub device: crate::signaling::Device,
pub presence_status: DevicePresenceStatus,
pub presence_updated_at: Option<i64>,
pub presence_expires_at: Option<i64>,
pub connectable: bool,
pub connection_status: DeviceConnectionStatus,
pub settled_ready: bool,
pub readiness_state: ReadinessState,
pub readiness_reason: String,
pub peer_health: crate::connection_manager::ConnectionHealth,
pub peer_id: Option<String>,
pub scopes: Vec<String>,
pub connection_id: Option<String>,
pub device_id_hint: Option<String>,
pub active_transport_stable_id: Option<u64>,
pub transport_generation: u64,
pub active_transport: String,
pub parallel_transport: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct PeerSessionSnapshot {
pub peer_id: String,
pub device_id: Option<String>,
pub device_id_hint: Option<String>,
pub node_id: Option<String>,
pub active_connection_id: Option<String>,
pub candidate_connection_ids: Vec<String>,
pub status: crate::connection_manager::ConnectionState,
pub health: crate::connection_manager::ConnectionHealth,
pub settled_ready: bool,
pub readiness_state: ReadinessState,
pub active_transport_stable_id: Option<u64>,
pub transport_generation: u64,
pub active_transport: String,
pub parallel_transport: Option<String>,
pub replacement_pending: bool,
pub last_lifecycle_transition_at_ms: i64,
pub readiness_reason: String,
pub transition_count: u64,
pub connecting_transition_count: u64,
pub replacement_count: u64,
pub retire_count: u64,
pub last_disconnect_reason: Option<String>,
pub last_reconnect_reason: Option<String>,
pub scopes: Vec<String>,
pub last_seen_at_ms: i64,
pub error: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ManagedConnectResult {
pub connection_id: String,
pub device_id: Option<String>,
pub device_id_hint: Option<String>,
pub remote_node_id: String,
pub state: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub approved_scope: Option<String>,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ManagedConnectionAdoption {
pub connection_id: String,
pub node_id: String,
pub device_id: Option<String>,
pub transport_generation: u64,
pub status_reason: Option<String>,
pub main_stream_ready: bool,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum ManagedConnectionHealthStatus {
Healthy,
AwaitingReplacement,
Dead,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ManagedConnectionHealthSnapshot {
pub connection_id: String,
pub device_id: Option<String>,
pub device_id_hint: Option<String>,
pub node_id: Option<String>,
pub active_transport_stable_id: Option<u64>,
pub transport_generation: u64,
pub status: ManagedConnectionHealthStatus,
pub settled_ready: bool,
pub readiness_state: ReadinessState,
pub replacement_pending: bool,
pub last_lifecycle_transition_at_ms: i64,
pub readiness_reason: String,
pub transition_count: u64,
pub connecting_transition_count: u64,
pub replacement_count: u64,
pub retire_count: u64,
pub last_disconnect_reason: Option<String>,
pub last_reconnect_reason: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ConnectionStateSnapshot {
pub connection_id: String,
pub device_id: Option<String>,
pub device_id_hint: Option<String>,
pub remote_node_id: Option<String>,
pub state: String,
pub transport_state: String,
pub protocol_state: String,
pub routable: bool,
pub readiness_state: ReadinessState,
pub readiness_reason: String,
pub transport_generation: u64,
pub active_transport_stable_id: Option<u64>,
pub active_transport: String,
pub parallel_transport: Option<String>,
pub replacement_in_progress: bool,
pub last_lifecycle_transition_at_ms: i64,
pub transition_count: u64,
pub connecting_transition_count: u64,
pub replacement_count: u64,
pub retire_count: u64,
pub last_disconnect_reason: Option<String>,
pub last_reconnect_reason: Option<String>,
pub error: Option<String>,
pub created_at: i64,
pub updated_at: i64,
}
#[cfg(target_arch = "wasm32")]
#[derive(Debug, Clone, PartialEq, Eq)]
struct WasmConnectionStateFingerprint {
connection_id: String,
device_id: Option<String>,
device_id_hint: Option<String>,
remote_node_id: Option<String>,
state: String,
transport_state: String,
protocol_state: String,
routable: bool,
transport_generation: u64,
active_transport: String,
parallel_transport: Option<String>,
replacement_in_progress: bool,
transition_count: u64,
connecting_transition_count: u64,
replacement_count: u64,
retire_count: u64,
last_disconnect_reason: Option<String>,
last_reconnect_reason: Option<String>,
error: Option<String>,
}
#[cfg(target_arch = "wasm32")]
impl From<&ConnectionStateSnapshot> for WasmConnectionStateFingerprint {
fn from(snapshot: &ConnectionStateSnapshot) -> Self {
Self {
connection_id: snapshot.connection_id.clone(),
device_id: snapshot.device_id.clone(),
device_id_hint: snapshot.device_id_hint.clone(),
remote_node_id: snapshot.remote_node_id.clone(),
state: snapshot.state.clone(),
transport_state: snapshot.transport_state.clone(),
protocol_state: snapshot.protocol_state.clone(),
routable: snapshot.routable,
transport_generation: snapshot.transport_generation,
active_transport: snapshot.active_transport.clone(),
parallel_transport: snapshot.parallel_transport.clone(),
replacement_in_progress: snapshot.replacement_in_progress,
transition_count: snapshot.transition_count,
connecting_transition_count: snapshot.connecting_transition_count,
replacement_count: snapshot.replacement_count,
retire_count: snapshot.retire_count,
last_disconnect_reason: snapshot.last_disconnect_reason.clone(),
last_reconnect_reason: snapshot.last_reconnect_reason.clone(),
error: snapshot.error.clone(),
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase", tag = "kind")]
pub enum ManagedConnectionBridgeAction {
Healthy,
AwaitReplacement {
reason: String,
},
Rebind {
remote_node_id: Option<String>,
transport_generation: u64,
},
Retire {
reason: String,
},
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DuplicateClosedHandling {
RebindKeptTransport,
RetireClosedTransport,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum IncomingTransportCloseResolution {
PreservedKeptTransport,
RetiredClosedTransport,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AutoConnectTieBreakDecision {
ObserveConnectedTransport,
WaitForInitiator,
ActAsInitiator,
}
#[cfg(not(target_arch = "wasm32"))]
fn duplicate_closed_handling(
close_reason_debug: &str,
kept_transport_alive: bool,
kept_transport_healthy: bool,
kept_stable_id: u64,
closed_stable_id: u64,
) -> DuplicateClosedHandling {
let _ = close_reason_debug;
let _ = kept_transport_healthy;
if kept_stable_id != closed_stable_id && kept_transport_alive {
DuplicateClosedHandling::RebindKeptTransport
} else {
DuplicateClosedHandling::RetireClosedTransport
}
}
#[cfg(not(target_arch = "wasm32"))]
fn incoming_local_close_should_wait_for_replacement(close_reason_debug: &str) -> bool {
if crate::lifecycle_reason::reason_is_replacement_churn(Some(close_reason_debug)) {
return true;
}
let normalized = close_reason_debug.trim().to_ascii_lowercase();
normalized.contains("timedout")
|| normalized.contains("timed out")
|| normalized.contains("timeout")
}
#[cfg(not(target_arch = "wasm32"))]
fn auto_connect_tie_break_decision(
local_node_id: &str,
remote_node_id: &str,
transport_connected: bool,
waited_ms: i64,
initiator_grace_ms: i64,
) -> AutoConnectTieBreakDecision {
if transport_connected {
return AutoConnectTieBreakDecision::ObserveConnectedTransport;
}
if local_node_id <= remote_node_id {
if waited_ms >= initiator_grace_ms {
AutoConnectTieBreakDecision::ActAsInitiator
} else {
AutoConnectTieBreakDecision::WaitForInitiator
}
} else {
AutoConnectTieBreakDecision::ActAsInitiator
}
}
fn normalize_lookup_id(value: Option<&str>) -> Option<String> {
let trimmed = value?.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_ascii_lowercase())
}
}
fn snapshot_connection_status(
device: &crate::signaling::Device,
peer: Option<&crate::connection_manager::PeerSnapshot>,
) -> DeviceConnectionStatus {
match peer.map(|value| &value.status) {
Some(crate::connection_manager::ConnectionState::Pending)
| Some(crate::connection_manager::ConnectionState::Connecting) => {
DeviceConnectionStatus::Connecting
}
Some(crate::connection_manager::ConnectionState::Connected)
if peer.map(peer_snapshot_settled_ready).unwrap_or(false) =>
{
DeviceConnectionStatus::Connected
}
Some(crate::connection_manager::ConnectionState::Connected) => {
DeviceConnectionStatus::Connecting
}
Some(crate::connection_manager::ConnectionState::Failed) => DeviceConnectionStatus::Failed,
Some(crate::connection_manager::ConnectionState::Closed)
| Some(crate::connection_manager::ConnectionState::Closing) => {
DeviceConnectionStatus::Closed
}
None if device.online => DeviceConnectionStatus::Online,
None => DeviceConnectionStatus::Disconnected,
}
}
fn snapshot_settled_ready(
peer: Option<&crate::connection_manager::PeerSnapshot>,
connection_status: &DeviceConnectionStatus,
) -> bool {
matches!(connection_status, DeviceConnectionStatus::Connected)
&& peer.map(peer_snapshot_settled_ready).unwrap_or(false)
}
fn device_presence_updated_at(device: &crate::signaling::Device) -> Option<i64> {
[
device.last_seen_at.as_ref(),
device.updated_at.as_ref(),
device.created_at.as_ref(),
]
.into_iter()
.flatten()
.filter_map(|value| crate::firebase::device_registry::parse_millis(Some(value)))
.max()
}
fn device_presence_expires_at(device: &crate::signaling::Device) -> Option<i64> {
crate::firebase::device_registry::parse_millis(device.expires_at.as_ref())
}
fn device_presence_status(
device: &crate::signaling::Device,
presence_updated_at: Option<i64>,
now_ms: i64,
) -> DevicePresenceStatus {
if !device.online {
return DevicePresenceStatus::Offline;
}
if presence_updated_at
.map(|updated_at| {
updated_at.saturating_add(crate::firebase::device_registry::DEVICE_STALE_HEARTBEAT_MS)
<= now_ms
})
.unwrap_or(false)
{
DevicePresenceStatus::Idle
} else {
DevicePresenceStatus::Online
}
}
fn device_connectable(
device: &crate::signaling::Device,
presence_status: &DevicePresenceStatus,
) -> bool {
matches!(presence_status, DevicePresenceStatus::Online)
&& device
.ticket
.as_deref()
.map(str::trim)
.is_some_and(|value| !value.is_empty())
}
fn peer_snapshot_settled_ready(peer: &crate::connection_manager::PeerSnapshot) -> bool {
matches!(peer_readiness_state(peer), ReadinessState::Routable)
}
fn peer_readiness_state(peer: &crate::connection_manager::PeerSnapshot) -> ReadinessState {
match peer.status {
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting => ReadinessState::Connecting,
crate::connection_manager::ConnectionState::Connected => {
let connection_ids_present = !peer.connection_ids.is_empty();
let healthy_with_connections = connection_ids_present
&& matches!(
peer.health,
crate::connection_manager::ConnectionHealth::Healthy
);
let health_explicitly_unhealthy = matches!(
peer.health,
crate::connection_manager::ConnectionHealth::Suspect
| crate::connection_manager::ConnectionHealth::Stale
);
let live_transport_with_connections = connection_ids_present
&& peer.active_transport_stable_id.is_some()
&& !health_explicitly_unhealthy;
if healthy_with_connections || live_transport_with_connections {
ReadinessState::Routable
} else if peer.active_transport_stable_id.is_none() {
ReadinessState::TransportOnly
} else if now_millis_i64().saturating_sub(peer.last_lifecycle_transition_at_ms)
>= MANAGED_SETTLE_DEADLINE_MS
{
ReadinessState::AwaitingReplacement
} else {
ReadinessState::Settling
}
}
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => ReadinessState::Closed,
crate::connection_manager::ConnectionState::Failed => ReadinessState::Failed,
}
}
fn peer_readiness_reason(peer: &crate::connection_manager::PeerSnapshot) -> String {
match peer_readiness_state(peer) {
ReadinessState::Connecting => "waiting-for-transport".to_string(),
ReadinessState::TransportOnly => {
"transport-attached-awaiting-readiness-confirmation".to_string()
}
ReadinessState::Settling => "transport-connected-health-check-pending".to_string(),
ReadinessState::Routable => "transport-healthy".to_string(),
ReadinessState::AwaitingReplacement => "readiness-confirmation-timed-out".to_string(),
ReadinessState::Closed => "transport-closed".to_string(),
ReadinessState::Failed => "connection-failed".to_string(),
}
}
fn connection_readiness_state(
record: &crate::connection_manager::ConnectionRecord,
peer_snapshot: Option<&crate::connection_manager::PeerSnapshot>,
) -> ReadinessState {
match record.state {
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting => ReadinessState::Connecting,
crate::connection_manager::ConnectionState::Connected => {
if peer_snapshot
.map(peer_snapshot_settled_ready)
.unwrap_or(false)
{
ReadinessState::Routable
} else if record.transport_stable_id.is_none() {
ReadinessState::AwaitingReplacement
} else if now_millis_i64().saturating_sub(record.last_transport_change_at_ms)
>= MANAGED_SETTLE_DEADLINE_MS
{
ReadinessState::AwaitingReplacement
} else {
ReadinessState::Settling
}
}
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => ReadinessState::Closed,
crate::connection_manager::ConnectionState::Failed => ReadinessState::Failed,
}
}
fn connection_readiness_reason(
record: &crate::connection_manager::ConnectionRecord,
peer_snapshot: Option<&crate::connection_manager::PeerSnapshot>,
) -> String {
match connection_readiness_state(record, peer_snapshot) {
ReadinessState::Connecting => "waiting-for-transport".to_string(),
ReadinessState::TransportOnly => {
"transport-attached-awaiting-readiness-confirmation".to_string()
}
ReadinessState::Settling => "bounded-health-confirmation-pending".to_string(),
ReadinessState::Routable => "transport-healthy".to_string(),
ReadinessState::AwaitingReplacement => {
if record.transport_stable_id.is_none() {
"missing-active-transport".to_string()
} else {
"readiness-confirmation-timed-out".to_string()
}
}
ReadinessState::Closed => "transport-closed".to_string(),
ReadinessState::Failed => "connection-failed".to_string(),
}
}
fn public_peer_state(
peer: &crate::connection_manager::PeerSnapshot,
) -> crate::connection_manager::ConnectionState {
if matches!(
peer.status,
crate::connection_manager::ConnectionState::Connected
) && !peer_snapshot_settled_ready(peer)
{
crate::connection_manager::ConnectionState::Connecting
} else {
peer.status.clone()
}
}
fn peer_session_snapshot_from_peer(
peer: crate::connection_manager::PeerSnapshot,
) -> PeerSessionSnapshot {
let status = public_peer_state(&peer);
let active_connection_id = peer.connection_ids.first().cloned();
let candidate_connection_ids = if peer.connection_ids.len() > 1 {
peer.connection_ids[1..].to_vec()
} else {
Vec::new()
};
let settled_ready = peer_snapshot_settled_ready(&peer);
let readiness_state = peer_readiness_state(&peer);
let readiness_reason = peer_readiness_reason(&peer);
let crate::connection_manager::PeerSnapshot {
peer_id,
device_id,
device_id_hint,
node_id,
active_transport_stable_id,
active_transport_generation,
active_transport,
parallel_transport,
last_lifecycle_transition_at_ms,
health,
transition_count,
connecting_transition_count,
replacement_count,
retire_count,
last_disconnect_reason,
last_reconnect_reason,
scopes,
last_seen_at_ms,
error,
..
} = peer;
PeerSessionSnapshot {
peer_id,
device_id,
device_id_hint,
node_id,
active_connection_id,
candidate_connection_ids,
status,
health,
settled_ready,
readiness_state: readiness_state.clone(),
active_transport_stable_id,
transport_generation: active_transport_generation,
active_transport,
parallel_transport,
replacement_pending: matches!(readiness_state, ReadinessState::AwaitingReplacement),
last_lifecycle_transition_at_ms,
readiness_reason,
transition_count,
connecting_transition_count,
replacement_count,
retire_count,
last_disconnect_reason,
last_reconnect_reason,
scopes,
last_seen_at_ms,
error,
}
}
fn connection_state_snapshot_from_parts(
record: &crate::connection_manager::ConnectionRecord,
peer_snapshot: Option<&crate::connection_manager::PeerSnapshot>,
) -> ConnectionStateSnapshot {
let settled_ready = peer_snapshot
.map(peer_snapshot_settled_ready)
.unwrap_or(false);
let readiness_state = connection_readiness_state(record, peer_snapshot);
let in_place_reconnect_close = matches!(
record.state,
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
) && crate::lifecycle_reason::LifecycleReasonCode::from_text(
record.status_reason.as_deref(),
)
.map(|code| {
matches!(
code,
crate::lifecycle_reason::LifecycleReasonCode::ReplacementInProgress
) || code.is_transient_reconnect()
})
.unwrap_or(false);
let state = match record.state {
crate::connection_manager::ConnectionState::Pending => "connecting",
crate::connection_manager::ConnectionState::Connecting => "connecting",
crate::connection_manager::ConnectionState::Connected if settled_ready => "connected",
crate::connection_manager::ConnectionState::Connected => "connecting",
crate::connection_manager::ConnectionState::Failed => "failed",
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
if in_place_reconnect_close =>
{
"connecting"
}
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => "closed",
};
let transport_state = match record.state {
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting => "connecting",
crate::connection_manager::ConnectionState::Connected => "connected",
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
if in_place_reconnect_close =>
{
"connecting"
}
crate::connection_manager::ConnectionState::Failed
| crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => "closed",
};
let protocol_state = match record.state {
crate::connection_manager::ConnectionState::Connected if settled_ready => "routable",
crate::connection_manager::ConnectionState::Connected => "transport-only",
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting => "connecting",
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
if in_place_reconnect_close =>
{
"connecting"
}
crate::connection_manager::ConnectionState::Failed
| crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => "closed",
};
ConnectionStateSnapshot {
connection_id: record.connection_id.clone(),
device_id: record
.device_id
.clone()
.or_else(|| peer_snapshot.and_then(|snapshot| snapshot.device_id.clone())),
device_id_hint: record
.device_id_hint
.clone()
.or_else(|| peer_snapshot.and_then(|snapshot| snapshot.device_id_hint.clone())),
remote_node_id: record.node_id.clone(),
state: state.to_string(),
transport_state: transport_state.to_string(),
protocol_state: protocol_state.to_string(),
routable: settled_ready,
readiness_state: readiness_state.clone(),
readiness_reason: connection_readiness_reason(record, peer_snapshot),
transport_generation: record.transport_generation,
active_transport_stable_id: record.transport_stable_id,
active_transport: record.active_transport.clone(),
parallel_transport: record.parallel_transport.clone(),
replacement_in_progress: matches!(readiness_state, ReadinessState::AwaitingReplacement),
last_lifecycle_transition_at_ms: record
.last_transport_change_at_ms
.max(record.updated_at_ms),
transition_count: record.transition_count,
connecting_transition_count: record.connecting_transition_count,
replacement_count: record.replacement_count,
retire_count: record.retire_count,
last_disconnect_reason: record.last_disconnect_reason.clone(),
last_reconnect_reason: record.last_reconnect_reason.clone(),
error: record.status_reason.clone(),
created_at: record.created_at_ms,
updated_at: record.updated_at_ms,
}
}
fn merge_device_status_snapshots(
devices: Vec<crate::signaling::Device>,
peers: Vec<crate::connection_manager::PeerSnapshot>,
) -> Vec<DeviceStatusSnapshot> {
let mut peer_by_device: std::collections::HashMap<
String,
crate::connection_manager::PeerSnapshot,
> = std::collections::HashMap::new();
for peer in peers {
let aliases = [
normalize_lookup_id(peer.device_id.as_deref()),
normalize_lookup_id(peer.device_id_hint.as_deref()),
normalize_lookup_id(peer.node_id.as_deref()),
normalize_lookup_id(Some(peer.peer_id.as_str())),
];
for alias in aliases.into_iter().flatten() {
let replace = match peer_by_device.get(&alias) {
Some(existing) => prefer_device_status_peer_snapshot(&peer, existing),
None => true,
};
if replace {
peer_by_device.insert(alias, peer.clone());
}
}
}
devices
.into_iter()
.map(|device| {
let peer = [
normalize_lookup_id(Some(device.device_id.as_str())),
normalize_lookup_id(device.node_id.as_deref()),
]
.into_iter()
.flatten()
.find_map(|key| peer_by_device.get(&key).cloned());
let connection_status = snapshot_connection_status(&device, peer.as_ref());
let presence_updated_at = device_presence_updated_at(&device);
let presence_expires_at = device_presence_expires_at(&device);
let presence_status =
device_presence_status(&device, presence_updated_at, now_millis_i64());
let connectable = device_connectable(&device, &presence_status);
let readiness_state = peer.as_ref().map(peer_readiness_state).unwrap_or_else(|| {
if device.online {
ReadinessState::Connecting
} else {
ReadinessState::Closed
}
});
DeviceStatusSnapshot {
presence_status,
presence_updated_at,
presence_expires_at,
connectable,
settled_ready: snapshot_settled_ready(peer.as_ref(), &connection_status),
readiness_state: readiness_state.clone(),
readiness_reason: peer.as_ref().map(peer_readiness_reason).unwrap_or_else(|| {
if device.online {
"device-online-awaiting-runtime-session".to_string()
} else {
"device-offline".to_string()
}
}),
connection_status,
peer_health: peer
.as_ref()
.map(|value| value.health.clone())
.unwrap_or_else(|| {
if device.online {
crate::connection_manager::ConnectionHealth::Healthy
} else {
crate::connection_manager::ConnectionHealth::Unknown
}
}),
peer_id: peer.as_ref().map(|value| value.peer_id.clone()),
scopes: peer
.as_ref()
.map(|value| value.scopes.clone())
.unwrap_or_default(),
connection_id: peer
.as_ref()
.and_then(|value| value.connection_ids.first().cloned()),
device_id_hint: peer.as_ref().and_then(|value| value.device_id_hint.clone()),
active_transport_stable_id: peer
.as_ref()
.and_then(|value| value.active_transport_stable_id),
transport_generation: peer
.as_ref()
.map(|value| value.active_transport_generation)
.unwrap_or(0),
active_transport: peer
.as_ref()
.map(|value| value.active_transport.clone())
.unwrap_or_else(|| "iroh".to_string()),
parallel_transport: peer
.as_ref()
.and_then(|value| value.parallel_transport.clone()),
device,
}
})
.collect()
}
fn peer_snapshot_recency_ms(peer: &crate::connection_manager::PeerSnapshot) -> i64 {
peer.last_seen_at_ms
.max(peer.last_lifecycle_transition_at_ms)
}
fn prefer_device_status_peer_snapshot(
candidate: &crate::connection_manager::PeerSnapshot,
existing: &crate::connection_manager::PeerSnapshot,
) -> bool {
use crate::connection_manager::ConnectionState;
let c = &candidate.status;
let e = &existing.status;
let c_fail = matches!(c, ConnectionState::Failed);
let c_conn = matches!(c, ConnectionState::Connected);
let e_fail = matches!(e, ConnectionState::Failed);
let e_conn = matches!(e, ConnectionState::Connected);
let c_live = c_conn
&& candidate.active_transport_stable_id.is_some()
&& !candidate.connection_ids.is_empty();
let e_live = e_conn
&& existing.active_transport_stable_id.is_some()
&& !existing.connection_ids.is_empty();
if c_live && !e_live {
return true;
}
if !c_live && e_live {
return false;
}
if c_fail && e_conn {
return peer_snapshot_recency_ms(candidate) >= peer_snapshot_recency_ms(existing);
}
if c_conn && e_fail {
return peer_snapshot_recency_ms(candidate) > peer_snapshot_recency_ms(existing);
}
device_status_peer_priority(candidate) > device_status_peer_priority(existing)
}
fn device_status_peer_priority(
peer: &crate::connection_manager::PeerSnapshot,
) -> (u8, u8, u8, u8, i64, u64, u64) {
let readiness_rank = match peer_readiness_state(peer) {
ReadinessState::Routable => 5,
ReadinessState::Settling => 4,
ReadinessState::TransportOnly => 3,
ReadinessState::Connecting => 2,
ReadinessState::AwaitingReplacement => 1,
ReadinessState::Closed | ReadinessState::Failed => 0,
};
let status_rank = match peer.status {
crate::connection_manager::ConnectionState::Connected => 3,
crate::connection_manager::ConnectionState::Connecting
| crate::connection_manager::ConnectionState::Pending => 2,
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed => 1,
crate::connection_manager::ConnectionState::Failed => 0,
};
let health_rank = match peer.health {
crate::connection_manager::ConnectionHealth::Healthy => 3,
crate::connection_manager::ConnectionHealth::Suspect => 2,
crate::connection_manager::ConnectionHealth::Unknown => 1,
crate::connection_manager::ConnectionHealth::Stale => 0,
};
(
readiness_rank,
status_rank,
health_rank,
u8::from(peer.active_transport_stable_id.is_some()),
peer.last_seen_at_ms,
peer.transition_count,
peer.replacement_count,
)
}
#[cfg(not(target_arch = "wasm32"))]
fn auto_connect_verbose() -> bool {
std::env::var("OPENRTC_AUTOCONNECT_VERBOSE")
.map(|value| {
let normalized = value.trim().to_ascii_lowercase();
normalized == "1" || normalized == "true" || normalized == "yes"
})
.unwrap_or(false)
}
#[cfg(not(target_arch = "wasm32"))]
fn parse_env_bool(value: &str) -> Option<bool> {
match value.trim().to_ascii_lowercase().as_str() {
"1" | "true" | "yes" | "on" => Some(true),
"0" | "false" | "no" | "off" => Some(false),
_ => None,
}
}
#[cfg(not(target_arch = "wasm32"))]
fn env_bool(name: &str) -> Option<bool> {
std::env::var(name)
.ok()
.and_then(|value| parse_env_bool(&value))
}
#[cfg(not(target_arch = "wasm32"))]
fn should_use_relay_only_mode() -> bool {
env_bool("PLUTO_IROH_RELAY_ONLY").unwrap_or(false)
}
#[cfg(not(target_arch = "wasm32"))]
fn should_disable_ipv6() -> bool {
if let Some(disable) = env_bool("PLUTO_IROH_DISABLE_IPV6") {
return disable;
}
if let Some(enable) = env_bool("PLUTO_IROH_ENABLE_IPV6") {
return !enable;
}
true
}
#[cfg(not(target_arch = "wasm32"))]
fn apply_native_network_preferences(
mut builder: iroh::endpoint::Builder,
relay_only: bool,
_relay_transport_policy: IrohRelayTransportPolicy,
) -> anyhow::Result<iroh::endpoint::Builder> {
{
use iroh::endpoint::{QuicTransportConfig, VarInt};
let idle_timeout_secs = if cfg!(target_arch = "wasm32") {
15
} else {
120
};
let idle_timeout: iroh::endpoint::IdleTimeout =
std::time::Duration::from_secs(idle_timeout_secs)
.try_into()
.map_err(|e| anyhow::anyhow!("invalid idle timeout: {}", e))?;
let stream_window = VarInt::from_u32(8 * 1024 * 1024); let conn_window = VarInt::from_u32(16 * 1024 * 1024); let send_window: u64 = 16 * 1024 * 1024; let keep_alive_secs: u64 = if cfg!(any(target_os = "ios", target_os = "android")) {
90
} else if cfg!(target_arch = "wasm32") {
5
} else {
25
};
let transport_config = QuicTransportConfig::builder()
.max_idle_timeout(Some(idle_timeout))
.keep_alive_interval(std::time::Duration::from_secs(keep_alive_secs))
.stream_receive_window(stream_window)
.receive_window(conn_window)
.send_window(send_window)
.build();
builder = builder.transport_config(transport_config);
}
if relay_only {
return Ok(builder.clear_ip_transports());
}
if should_disable_ipv6() {
builder = builder.clear_ip_transports();
builder = builder
.bind_addr("0.0.0.0:0")
.map_err(|e| anyhow::anyhow!("failed to bind IPv4 transport: {}", e))?;
}
Ok(builder)
}
#[cfg(not(target_arch = "wasm32"))]
fn auto_connect_failure_backoff_ms(failure_count: u8) -> i64 {
match failure_count {
0 => 0,
_ => {
let exp = 1_000_i64.saturating_mul(1_i64 << (failure_count as u32).min(5));
exp.min(30_000)
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn auto_connect_network_change_failure_threshold() -> u8 {
std::env::var("OPENRTC_NETWORK_CHANGE_FAILURE_THRESHOLD")
.ok()
.and_then(|value| value.parse::<u8>().ok())
.map(|value| value.clamp(1, 6))
.unwrap_or(2)
}
#[cfg(not(target_arch = "wasm32"))]
fn non_initiator_escalation_grace_ms(escalation_count: u8) -> i64 {
let base_ms = 2_000_i64;
let exp_ms = base_ms.saturating_mul(1_i64 << (escalation_count as u32).min(3));
let capped_ms = exp_ms.min(16_000);
let jitter_pct = ((escalation_count as i64 * 37 + 7) % 41) - 20; let jitter_ms = capped_ms * jitter_pct / 100;
(capped_ms + jitter_ms).max(1_000)
}
#[cfg(not(target_arch = "wasm32"))]
fn register_auto_connect_failure(
remote_device_id: &str,
now_ms: i64,
failure_count: &mut std::collections::HashMap<String, u8>,
failure_backoff_until: &mut std::collections::HashMap<String, i64>,
) {
let count = failure_count
.entry(remote_device_id.to_string())
.or_insert(0);
*count = count.saturating_add(1).min(10);
let backoff_ms = auto_connect_failure_backoff_ms(*count);
failure_backoff_until.insert(
remote_device_id.to_string(),
now_ms.saturating_add(backoff_ms),
);
}
#[cfg(not(target_arch = "wasm32"))]
fn clear_auto_connect_failure_state(
remote_device_id: &str,
failure_count: &mut std::collections::HashMap<String, u8>,
failure_backoff_until: &mut std::collections::HashMap<String, i64>,
) {
failure_count.remove(remote_device_id);
failure_backoff_until.remove(remote_device_id);
}
#[cfg(target_arch = "wasm32")]
fn emit_wasm_connection_state_event(snapshot: &ConnectionStateSnapshot) {
use wasm_bindgen::JsValue;
let Some(window) = web_sys::window() else {
return;
};
let detail = js_sys::Object::new();
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("connectionId"),
&JsValue::from_str(snapshot.connection_id.as_str()),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("state"),
&JsValue::from_str(snapshot.state.as_str()),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("createdAt"),
&JsValue::from_f64(snapshot.created_at as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("updatedAt"),
&JsValue::from_f64(snapshot.updated_at as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("transportState"),
&JsValue::from_str(snapshot.transport_state.as_str()),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("protocolState"),
&JsValue::from_str(snapshot.protocol_state.as_str()),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("routable"),
&JsValue::from_bool(snapshot.routable),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("transportGeneration"),
&JsValue::from_f64(snapshot.transport_generation as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("activeTransport"),
&JsValue::from_str(snapshot.active_transport.as_str()),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("parallelTransport"),
&snapshot
.parallel_transport
.as_deref()
.map(JsValue::from_str)
.unwrap_or(JsValue::NULL),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("replacementInProgress"),
&JsValue::from_bool(snapshot.replacement_in_progress),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("transitionCount"),
&JsValue::from_f64(snapshot.transition_count as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("connectingTransitionCount"),
&JsValue::from_f64(snapshot.connecting_transition_count as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("replacementCount"),
&JsValue::from_f64(snapshot.replacement_count as f64),
);
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("retireCount"),
&JsValue::from_f64(snapshot.retire_count as f64),
);
if let Some(device_id) = snapshot.device_id.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("deviceId"),
&JsValue::from_str(device_id),
);
}
if let Some(device_id_hint) = snapshot.device_id_hint.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("deviceIdHint"),
&JsValue::from_str(device_id_hint),
);
}
if let Some(remote_node_id) = snapshot.remote_node_id.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("remoteNodeId"),
&JsValue::from_str(remote_node_id),
);
}
if let Some(error) = snapshot.error.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("error"),
&JsValue::from_str(error),
);
}
if let Some(reason) = snapshot.last_disconnect_reason.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("lastDisconnectReason"),
&JsValue::from_str(reason),
);
}
if let Some(reason) = snapshot.last_reconnect_reason.as_deref() {
let _ = js_sys::Reflect::set(
&detail,
&JsValue::from_str("lastReconnectReason"),
&JsValue::from_str(reason),
);
}
let init = web_sys::CustomEventInit::new();
init.set_detail(&detail.into());
if let Ok(event) =
web_sys::CustomEvent::new_with_event_init_dict("connection-state-changed", &init)
{
let _ = window.dispatch_event(&event);
}
}
fn parse_endpoint_ticket(ticket: &str) -> anyhow::Result<iroh::EndpointAddr> {
let trimmed = ticket.trim();
if trimmed.is_empty() {
return Err(anyhow::anyhow!("endpoint ticket is required"));
}
let parsed = EndpointTicket::from_str(trimmed)
.map_err(|e| anyhow::anyhow!("invalid endpoint ticket: {}", e))?;
Ok(parsed.endpoint_addr().clone())
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RuntimeStatus {
pub ready: bool,
pub node_id: Option<String>,
pub transport: RuntimeTransportStatus,
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RuntimeTransportFeatureStatus {
pub compiled: bool,
pub enabled: bool,
}
impl RuntimeTransportFeatureStatus {
pub(crate) fn new(compiled: bool, enabled: bool) -> Self {
Self {
compiled,
enabled: compiled && enabled,
}
}
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RuntimeTransportStatus {
pub iroh_quic: RuntimeTransportFeatureStatus,
pub iroh_lan: RuntimeTransportFeatureStatus,
pub web_rtc: RuntimeTransportFeatureStatus,
pub web_rtc_lan: RuntimeTransportFeatureStatus,
pub moq: RuntimeTransportFeatureStatus,
pub ble: RuntimeTransportFeatureStatus,
pub iroh_relay_only: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub iroh_relay_transport_policy: Option<IrohRelayTransportPolicy>,
}
impl RuntimeTransportStatus {
pub(crate) fn from_config(config: &TransportConfig) -> Self {
let lan_compiled = cfg!(feature = "transport-lan");
let lan_enabled = config
.iroh_lan
.as_ref()
.map(|lan| lan.enabled)
.unwrap_or(false);
let webrtc_compiled = cfg!(feature = "transport-webrtc");
let webrtc_enabled = config.webrtc.is_some();
let webrtc_lan_enabled = config
.webrtc
.as_ref()
.map(|webrtc| webrtc.lan_mode)
.unwrap_or(false);
let moq_compiled = cfg!(feature = "transport-moq");
let moq_enabled = config.moq.is_some();
let ble_compiled = cfg!(all(not(target_arch = "wasm32"), openrtc_unpublished_ble));
let ble_enabled = config.ble.as_ref().map(|ble| ble.enabled).unwrap_or(false);
Self {
iroh_quic: RuntimeTransportFeatureStatus::new(true, true),
iroh_lan: RuntimeTransportFeatureStatus::new(lan_compiled, lan_enabled),
web_rtc: RuntimeTransportFeatureStatus::new(webrtc_compiled, webrtc_enabled),
web_rtc_lan: RuntimeTransportFeatureStatus::new(
webrtc_compiled,
webrtc_enabled && webrtc_lan_enabled,
),
moq: RuntimeTransportFeatureStatus::new(moq_compiled, moq_enabled),
ble: RuntimeTransportFeatureStatus::new(ble_compiled, ble_enabled),
iroh_relay_only: config.iroh_relay_only,
iroh_relay_transport_policy: config.iroh_relay_transport_policy.clone(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub enum AuthMode {
External,
Anonymous,
Required,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct IceServerConfig {
#[serde(default)]
pub urls: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub username: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub credential: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct WebRTCConfig {
#[serde(default)]
pub ice_servers: Vec<IceServerConfig>,
#[serde(default)]
pub privacy_mode: bool,
#[serde(default)]
pub lan_mode: bool,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct MoQConfig {
#[serde(default = "default_moq_relay_url")]
pub relay_url: String,
}
fn default_moq_relay_url() -> String {
"https://relay.cloudflare.mediaoverquic.com:443/moq".to_string()
}
impl Default for MoQConfig {
fn default() -> Self {
Self {
relay_url: default_moq_relay_url(),
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum IrohRelayTransportPolicy {
Auto,
QuicRequired,
WebsocketRequired,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct TransportConfig {
#[serde(default)]
pub iroh_relay_only: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub iroh_relay_transport_policy: Option<IrohRelayTransportPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub iroh_lan: Option<IrohLanConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub webrtc: Option<WebRTCConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub moq: Option<MoQConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ble: Option<BleConfig>,
}
impl Default for TransportConfig {
fn default() -> Self {
Self {
iroh_relay_only: false,
iroh_relay_transport_policy: None,
iroh_lan: Some(IrohLanConfig::default()),
webrtc: None,
moq: None,
ble: None,
}
}
}
pub struct ClientBuilder {
project_id: String,
app_tag: String,
auth_mode: AuthMode,
transport_config: TransportConfig,
token_provider: Arc<dyn Fn() -> Option<String> + Send + Sync>,
signaling: Option<Arc<dyn SignalingBackend>>,
room: Option<Arc<dyn RoomBackend>>,
}
impl ClientBuilder {
pub fn new(
project_id: String,
api_key: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
Self::new_with_app_tag(
project_id,
crate::app_tag_from_api_key(&api_key),
token_provider,
)
}
pub fn new_with_app_tag(
project_id: String,
app_tag: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
Self {
project_id,
app_tag,
auth_mode: AuthMode::External,
transport_config: TransportConfig::default(),
token_provider: Arc::from(token_provider),
signaling: None,
room: None,
}
}
pub fn auth_mode(mut self, auth_mode: AuthMode) -> Self {
self.auth_mode = auth_mode;
self
}
pub fn transport_config(mut self, transport_config: TransportConfig) -> Self {
self.transport_config = transport_config;
self
}
pub fn signaling_backend(mut self, signaling: Arc<dyn SignalingBackend>) -> Self {
self.signaling = Some(signaling);
self
}
pub fn room_backend(mut self, room: Arc<dyn RoomBackend>) -> Self {
self.room = Some(room);
self
}
pub fn build(self) -> Client {
let firestore = FirestoreClient::new(self.project_id.clone(), {
let token_provider = self.token_provider.clone();
Box::new(move || (token_provider)())
});
let app_backgrounded = Arc::new(AtomicBool::new(false));
let room = self.room.unwrap_or_else(|| {
Arc::new(FirebaseRoomBackend::new(
firestore.clone(),
self.app_tag.clone(),
))
});
let signaling = self.signaling.unwrap_or_else(|| {
#[cfg(target_arch = "wasm32")]
{
Arc::new(SignalingBackendImpl::new(firestore, self.app_tag.clone()))
as Arc<dyn SignalingBackend>
}
#[cfg(any(target_os = "ios", target_os = "android"))]
{
let app_backgrounded_provider = {
let app_backgrounded = app_backgrounded.clone();
Arc::new(move || app_backgrounded.load(Ordering::Relaxed))
};
Arc::new(FirebaseSignalingBackend::new_with_background_provider(
firestore,
self.app_tag.clone(),
app_backgrounded_provider,
)) as Arc<dyn SignalingBackend>
}
#[cfg(all(
not(target_arch = "wasm32"),
not(any(target_os = "ios", target_os = "android"))
))]
{
let app_backgrounded_provider = {
let app_backgrounded = app_backgrounded.clone();
Arc::new(move || app_backgrounded.load(Ordering::Relaxed))
};
Arc::new(SignalingBackendImpl::new(
&self.project_id,
self.app_tag.clone(),
self.token_provider.clone(),
app_backgrounded_provider,
)) as Arc<dyn SignalingBackend>
}
});
#[cfg(not(target_arch = "wasm32"))]
let (native_device_updates, _) = tokio::sync::broadcast::channel(32);
#[cfg(not(target_arch = "wasm32"))]
let (native_connection_state_updates, _) = tokio::sync::broadcast::channel(64);
#[cfg(not(target_arch = "wasm32"))]
let (native_peer_data_updates, _) = tokio::sync::broadcast::channel(256);
Client {
app_tag: self.app_tag,
signaling,
room,
node_id: Arc::new(RwLock::new(None)),
iroh_endpoint: Arc::new(RwLock::new(None)),
#[cfg(target_arch = "wasm32")]
iroh_node: Arc::new(RwLock::new(None)),
#[cfg(not(target_arch = "wasm32"))]
iroh_node: Arc::new(RwLock::new(None)),
connection_manager: Arc::new(crate::connection_manager::ConnectionManager::new()),
session_token_registry: Arc::new(crate::session_token::SessionTokenRegistry::new()),
connection_application_crypto_keys:
crate::client::application_crypto_impl::new_connection_application_crypto_key_map(),
connection_application_crypto_required:
crate::client::application_crypto_impl::new_connection_application_crypto_required_set(),
connection_application_crypto_outbound_sequences:
crate::client::application_crypto_impl::new_connection_application_crypto_outbound_sequences(),
connection_application_key_agreements:
crate::client::application_crypto_impl::new_connection_application_key_agreement_map(),
#[cfg(not(target_arch = "wasm32"))]
managed_scope_tickets: Arc::new(StdRwLock::new(HashMap::new())),
#[cfg(not(target_arch = "wasm32"))]
native_device_identity: Arc::new(RwLock::new(None)),
#[cfg(not(target_arch = "wasm32"))]
native_device_base_dir: Arc::new(RwLock::new(None)),
#[cfg(not(target_arch = "wasm32"))]
native_device_updates,
#[cfg(not(target_arch = "wasm32"))]
native_connection_state_updates,
#[cfg(not(target_arch = "wasm32"))]
native_peer_data_updates,
auto_connect_loop_key: Arc::new(Mutex::new(None)),
auto_connect_generation: Arc::new(AtomicU64::new(0)),
auto_connect_excluded: Arc::new(Mutex::new(HashSet::new())),
auto_connect_excluded_node_aliases: Arc::new(Mutex::new(HashMap::new())),
app_backgrounded,
#[cfg(target_arch = "wasm32")]
wasm_accept_bridge_started: Arc::new(std::sync::atomic::AtomicBool::new(false)),
#[cfg(target_arch = "wasm32")]
last_emitted_connection_states: Arc::new(Mutex::new(std::collections::HashMap::new())),
#[cfg(target_arch = "wasm32")]
last_empty_peer_sessions_warning_ms: Arc::new(AtomicU64::new(0)),
iroh_init_guard: Arc::new(tokio::sync::Mutex::new(())),
managed_connect_gates: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
presence_loop_tx: Arc::new(Mutex::new(None)),
#[cfg(not(target_arch = "wasm32"))]
app_limits: Arc::new(RwLock::new(AppLimits::default())),
#[cfg(not(target_arch = "wasm32"))]
project_id: self.project_id,
#[cfg(not(target_arch = "wasm32"))]
token_provider: self.token_provider,
transport_config: Arc::new(RwLock::new(self.transport_config)),
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "transport-lan", openrtc_unpublished_ble)
))]
local_discovery_registry: crate::local_discovery::LocalDiscoveryRegistry::new(),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
mdns_address_lookup: Arc::new(RwLock::new(None)),
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
ble_transport: Arc::new(RwLock::new(None)),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_sessions: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_suppressions: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_attempt_counts: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_start_gates: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_attempts_in_flight: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_webrtc_retry_deadlines: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
deferred_managed_retirements: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
native_moq_sessions: Arc::new(RwLock::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
native_signal_streams: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
scoped_connection_actor_registry: Arc::new(RwLock::new(None)),
#[cfg(not(target_arch = "wasm32"))]
auth_readiness: Arc::new(crate::client::auth_readiness::AuthReadinessStore::new()),
#[cfg(not(target_arch = "wasm32"))]
scope_classifier: Arc::new(RwLock::new(default_scope_classifier())),
}
}
}
mod admission_impl;
mod application_crypto_impl;
#[cfg(not(target_arch = "wasm32"))]
pub mod auth_readiness;
mod auto_connect_impl;
mod core_impl;
#[cfg(not(target_arch = "wasm32"))]
pub mod correlation;
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
mod drive_grant_actor;
#[cfg(all(
test,
not(target_arch = "wasm32"),
feature = "experimental-scoped-actor"
))]
mod drive_grant_actor_tests;
#[cfg(not(target_arch = "wasm32"))]
pub mod scope_classifier;
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
pub mod scoped_connection_actor;
mod state_signaling_impl;
#[cfg(not(target_arch = "wasm32"))]
mod transport_upgrade_impl;
#[cfg(target_arch = "wasm32")]
impl Client {
pub async fn send_peer(&self, _id: &str, _data: &[u8]) -> anyhow::Result<()> {
Err(anyhow::anyhow!(
"send_peer: use TypeScript Connection.sendTyped() for transport-aware sending in browser environments"
))
}
}
#[cfg(not(target_arch = "wasm32"))]
impl Client {
pub fn subscribe_native_peer_data(
&self,
) -> tokio::sync::broadcast::Receiver<NativePeerDataEvent> {
self.native_peer_data_updates.subscribe()
}
#[cfg_attr(
not(any(feature = "transport-webrtc", feature = "transport-moq")),
allow(dead_code)
)]
pub(crate) fn emit_native_peer_data(&self, event: NativePeerDataEvent) {
let _ = self.native_peer_data_updates.send(event);
}
}
#[cfg(test)]
mod tests;