use super::*;
pub(crate) fn should_represent_session_token_on_connected_transport(
has_extracted_token: bool,
session_admission_blocked: bool,
) -> bool {
has_extracted_token && session_admission_blocked
}
fn snapshot_has_independent_transport(snapshot: &crate::connection_manager::PeerSnapshot) -> bool {
crate::transport_label::is_independent_transport(snapshot.active_transport.as_str())
|| snapshot
.parallel_transport
.as_deref()
.map(crate::transport_label::is_independent_transport)
.unwrap_or(false)
}
#[cfg(target_arch = "wasm32")]
pub(crate) fn should_delay_wasm_stale_settled_snapshot_retirement(
last_transition_at_ms: i64,
now_ms: i64,
) -> bool {
now_ms.saturating_sub(last_transition_at_ms) < 4_500
}
impl Client {
async fn managed_connect_gate(&self, connection_id: &str) -> Arc<tokio::sync::Mutex<()>> {
let mut gates = self.managed_connect_gates.lock().await;
gates
.entry(connection_id.to_string())
.or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
.clone()
}
pub(crate) async fn retire_managed_connection_now(
&self,
connection_id: &str,
normalized_reason: Option<String>,
) {
self.connection_manager
.set_closing(connection_id, normalized_reason.clone())
.await;
self.connection_manager
.set_closed(connection_id, normalized_reason)
.await;
#[cfg(target_arch = "wasm32")]
self.emit_current_wasm_connection_state(connection_id).await;
let _ = self.connection_manager.remove(connection_id).await;
self.managed_connect_gates
.lock()
.await
.remove(connection_id);
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
if let Some(session) = self
.native_webrtc_sessions
.write()
.await
.remove(connection_id)
{
session.close();
}
self.deferred_managed_retirements
.write()
.await
.remove(connection_id);
self.native_signal_streams
.lock()
.await
.remove(connection_id);
self.native_webrtc_suppressions
.write()
.await
.remove(connection_id);
self.native_webrtc_attempt_counts
.write()
.await
.remove(connection_id);
self.native_webrtc_attempts_in_flight
.write()
.await
.remove(connection_id);
self.native_webrtc_retry_deadlines
.lock()
.await
.remove(connection_id);
if let Some(gate) = self
.native_webrtc_start_gates
.lock()
.await
.remove(connection_id)
{
gate.notify_waiters();
}
}
}
#[cfg(not(target_arch = "wasm32"))]
const SETTLED_PEER_STABILITY_WINDOW_MS: u64 = 3_000;
#[cfg(target_arch = "wasm32")]
const SETTLED_PEER_STABILITY_WINDOW_MS: u64 = 3_000;
#[cfg(target_arch = "wasm32")]
const EMPTY_PEER_SESSIONS_WARNING_INTERVAL_MS: u64 = 5_000;
fn normalize_transport_name(value: &str) -> Option<String> {
crate::transport_label::normalize(value).map(str::to_string)
}
pub async fn report_transport_status(
&self,
connection_id: &str,
active_transport: &str,
parallel_transport: Option<&str>,
) -> Option<PeerSessionSnapshot> {
let normalized_active_transport = Self::normalize_transport_name(active_transport)?;
let normalized_parallel_transport = parallel_transport
.and_then(Self::normalize_transport_name)
.filter(|value| value != &normalized_active_transport);
let current_record = self
.connection_manager
.get_by_connection_id(connection_id)
.await;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
let (active_transport, parallel_transport) = {
let mut active_transport = normalized_active_transport;
let mut parallel_transport = normalized_parallel_transport;
if let Some(record) = current_record.as_ref() {
if crate::transport_label::is_iroh_base(active_transport.as_str())
&& crate::native_moq_policy::should_preserve_active_moq_on_iroh_send(
record.active_transport.as_str(),
self.native_moq_data_ready_for_peer(connection_id).await,
)
{
active_transport = record.active_transport.clone();
parallel_transport = record.parallel_transport.clone();
}
}
(active_transport, parallel_transport)
};
#[cfg(any(target_arch = "wasm32", not(feature = "transport-moq")))]
let (active_transport, parallel_transport) =
(normalized_active_transport, normalized_parallel_transport);
let transport_changed = current_record
.as_ref()
.map(|record| {
record.active_transport != active_transport
|| record.parallel_transport != parallel_transport
})
.unwrap_or(true);
let snapshot = self
.connection_manager
.report_transport_status(connection_id, active_transport, parallel_transport)
.await
.map(peer_session_snapshot_from_peer)
.map(|snapshot| self.apply_peer_session_admission_guard(snapshot));
if snapshot.is_some() && transport_changed {
#[cfg(not(target_arch = "wasm32"))]
self.emit_current_native_connection_state(connection_id)
.await;
#[cfg(target_arch = "wasm32")]
self.emit_current_wasm_connection_state(connection_id).await;
}
snapshot
}
async fn live_transport_lookup(&self, id: &str) -> Option<(String, iroh::EndpointId, String)> {
let trimmed_id = id.trim();
if trimmed_id.is_empty() {
return None;
}
let local_node_id = self.current_node_id().await?;
let mut remote_node_candidates = Vec::new();
if let Ok(endpoint_id) = trimmed_id.parse::<iroh::EndpointId>() {
remote_node_candidates.push(endpoint_id.to_string());
}
if let Some((left, right)) = trimmed_id.split_once('-') {
if left == local_node_id {
remote_node_candidates.push(right.to_string());
} else if right == local_node_id {
remote_node_candidates.push(left.to_string());
}
}
if let Some(record) = self
.connection_manager
.get_by_connection_id(trimmed_id)
.await
{
if let Some(node_id) = record.node_id.as_deref() {
remote_node_candidates.push(node_id.to_string());
}
}
if let Some(record) = self
.connection_manager
.get_by_node_id(trimmed_id)
.await
.into_iter()
.max_by(|left, right| left.updated_at_ms.cmp(&right.updated_at_ms))
{
if let Some(node_id) = record.node_id.as_deref() {
remote_node_candidates.push(node_id.to_string());
}
}
if let Some(peer_snapshot) = self.connection_manager.peer_snapshot(trimmed_id).await {
if let Some(node_id) = peer_snapshot.node_id.as_deref() {
remote_node_candidates.push(node_id.to_string());
}
}
if let Some((user_id, _)) = self.active_session_identity() {
if let Ok(devices) = self.search_devices(&user_id).await {
if let Some(device) = devices
.iter()
.find(|device| device.device_id.trim() == trimmed_id)
{
if let Some(node_id) = device
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
remote_node_candidates.push(node_id.to_string());
}
}
}
}
let mut seen = std::collections::HashSet::new();
for remote_node_id in remote_node_candidates {
if remote_node_id.is_empty() || !seen.insert(remote_node_id.clone()) {
continue;
}
let Ok(endpoint_id) = remote_node_id.parse::<iroh::EndpointId>() else {
continue;
};
if !self.is_connection_transport_alive(endpoint_id).await {
continue;
}
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
return Some((connection_id, endpoint_id, remote_node_id));
}
None
}
async fn adopt_live_transport_session(
&self,
id: &str,
source: &str,
) -> Option<PeerSessionSnapshot> {
let (connection_id, endpoint_id, remote_node_id) = self.live_transport_lookup(id).await?;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64)?;
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
None,
Some(remote_node_id.clone()),
)
.await;
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.clone()),
Some(transport_stable_id),
Some(source.to_string()),
)
.await;
if let Some(device_id) = self
.known_remote_device_id_for_incoming_transport(&connection_id, &remote_node_id)
.await
{
self.mark_trusted_user_device_connection_admitted(&connection_id, &device_id)
.await;
self.prune_stale_records_for_device_node(&device_id, &remote_node_id)
.await;
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
self.connection_manager
.peer_snapshot(&connection_id)
.await
.map(peer_session_snapshot_from_peer)
}
async fn refresh_cached_peer_snapshot(
&self,
id: &str,
) -> Option<crate::connection_manager::PeerSnapshot> {
let snapshot = self.connection_manager.peer_snapshot(id).await?;
let endpoint_id = snapshot
.node_id
.as_deref()
.and_then(|value| value.parse::<iroh::EndpointId>().ok());
let transport_alive = match endpoint_id {
Some(endpoint_id) => Some(self.is_connection_transport_alive(endpoint_id).await),
None => None,
};
let live_transport_stable_id = match (endpoint_id, transport_alive) {
(Some(endpoint_id), Some(true)) => self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64),
_ => None,
};
let should_rebind_live_transport = matches!(transport_alive, Some(true))
&& (!matches!(
snapshot.status,
crate::connection_manager::ConnectionState::Connected
) || snapshot.active_transport_stable_id != live_transport_stable_id);
if should_rebind_live_transport {
let _ = self
.adopt_live_transport_session(id, "peer-snapshot-live-rebind")
.await;
return self
.connection_manager
.peer_snapshot(id)
.await
.or(Some(snapshot));
}
let should_retire = matches!(
snapshot.status,
crate::connection_manager::ConnectionState::Connected
) && peer_snapshot_settled_ready(&snapshot)
&& !snapshot_has_independent_transport(&snapshot)
&& matches!(transport_alive, Some(false));
if !should_retire {
#[cfg(target_arch = "wasm32")]
if peer_snapshot_settled_ready(&snapshot)
&& snapshot_has_independent_transport(&snapshot)
&& matches!(transport_alive, Some(false))
{
web_sys::console::debug_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][peer-session][wasm] preserving settled upgraded snapshot peer_id={} active_transport={} parallel_transport={:?} node_id={:?}",
snapshot.peer_id,
snapshot.active_transport,
snapshot.parallel_transport,
snapshot.node_id
)));
}
return Some(snapshot);
}
#[cfg(not(target_arch = "wasm32"))]
{
return Some(snapshot);
}
#[cfg(target_arch = "wasm32")]
if should_delay_wasm_stale_settled_snapshot_retirement(
snapshot.last_lifecycle_transition_at_ms,
now_millis_i64(),
) {
let mut awaiting_replacement = snapshot.clone();
awaiting_replacement.health = crate::connection_manager::ConnectionHealth::Stale;
awaiting_replacement.error =
Some("stale-settled-peer-snapshot-awaiting-replacement".to_string());
web_sys::console::debug_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][peer-session][wasm] delaying stale settled snapshot retirement during replacement grace but demoting readiness peer_id={} node_id={:?} connection_ids={:?}",
awaiting_replacement.peer_id,
awaiting_replacement.node_id,
awaiting_replacement.connection_ids
)));
return Some(awaiting_replacement);
}
#[cfg(target_arch = "wasm32")]
web_sys::console::debug_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][peer-session][wasm] retiring stale settled snapshot peer_id={} node_id={:?} connection_ids={:?}",
snapshot.peer_id,
snapshot.node_id,
snapshot.connection_ids
)));
#[cfg(target_arch = "wasm32")]
{
for connection_id in snapshot.connection_ids.clone() {
self.forget_session_connection(&connection_id);
self.retire_managed_connection(
&connection_id,
Some("stale-settled-peer-snapshot".to_string()),
)
.await;
}
self.connection_manager.peer_snapshot(id).await
}
}
fn session_admission_block_reason(&self, connection_id: &str) -> Option<(bool, String)> {
if !self.session_registry_active() {
return None;
}
match self.session_admission(connection_id) {
crate::session_token::SessionAdmission::Accepted { .. } => None,
crate::session_token::SessionAdmission::Pending => {
Some((false, "session-admission-pending".to_string()))
}
crate::session_token::SessionAdmission::Rejected { reason } => Some((true, reason)),
}
}
fn apply_peer_session_admission_guard(
&self,
mut snapshot: PeerSessionSnapshot,
) -> PeerSessionSnapshot {
let connection_id = snapshot
.active_connection_id
.clone()
.or_else(|| snapshot.candidate_connection_ids.first().cloned());
let Some(connection_id) = connection_id else {
return snapshot;
};
let Some((rejected, reason)) = self.session_admission_block_reason(&connection_id) else {
return snapshot;
};
snapshot.settled_ready = false;
snapshot.replacement_pending = false;
snapshot.readiness_reason = reason.clone();
if rejected {
snapshot.status = crate::connection_manager::ConnectionState::Failed;
snapshot.readiness_state = ReadinessState::Failed;
snapshot.health = crate::connection_manager::ConnectionHealth::Stale;
snapshot.error = Some(reason);
} else {
if matches!(
snapshot.status,
crate::connection_manager::ConnectionState::Connected
) {
snapshot.status = crate::connection_manager::ConnectionState::Connecting;
}
snapshot.readiness_state = ReadinessState::Settling;
}
snapshot
}
fn apply_connection_state_admission_guard(
&self,
mut snapshot: ConnectionStateSnapshot,
) -> ConnectionStateSnapshot {
let Some((rejected, reason)) = self.session_admission_block_reason(&snapshot.connection_id)
else {
return snapshot;
};
snapshot.routable = false;
snapshot.replacement_in_progress = false;
snapshot.readiness_reason = reason.clone();
if rejected {
snapshot.state = "failed".to_string();
snapshot.protocol_state = "admission-rejected".to_string();
snapshot.readiness_state = ReadinessState::Failed;
snapshot.error = Some(reason);
} else {
snapshot.state = "connecting".to_string();
snapshot.protocol_state = "awaiting-admission".to_string();
snapshot.readiness_state = ReadinessState::Settling;
}
snapshot
}
fn apply_device_status_admission_guard(
&self,
mut snapshot: DeviceStatusSnapshot,
) -> DeviceStatusSnapshot {
let Some(connection_id) = snapshot.connection_id.clone() else {
return snapshot;
};
let Some((rejected, reason)) = self.session_admission_block_reason(&connection_id) else {
return snapshot;
};
snapshot.settled_ready = false;
snapshot.readiness_reason = reason.clone();
if rejected {
snapshot.connection_status = DeviceConnectionStatus::Failed;
snapshot.readiness_state = ReadinessState::Failed;
snapshot.peer_health = crate::connection_manager::ConnectionHealth::Stale;
} else {
snapshot.connection_status = DeviceConnectionStatus::Connecting;
snapshot.readiness_state = ReadinessState::Settling;
}
snapshot
}
async fn resolve_settled_peer_endpoint(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(PeerSessionSnapshot, iroh::EndpointId)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self
.wait_for_settled_peer(id, Some(timeout_ms))
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
if !snapshot.settled_ready {
return Err(anyhow::anyhow!(
"peer {} is not settled for protocol traffic",
id
));
}
let node_id = snapshot
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("peer {} is missing a remote node id", id))?;
let endpoint_id = node_id.parse::<iroh::EndpointId>().map_err(|error| {
anyhow::anyhow!("peer {} has invalid remote node id: {}", id, error)
})?;
if self.is_connection_transport_alive(endpoint_id).await {
return Ok((snapshot, endpoint_id));
}
if let Some(connection_id) = snapshot.active_connection_id.as_deref() {
let _ = self
.confirm_managed_connection_readiness(connection_id)
.await;
}
#[cfg(not(target_arch = "wasm32"))]
let timed_out = started.elapsed().as_millis() >= timeout_ms as u128;
#[cfg(target_arch = "wasm32")]
let timed_out = (js_sys::Date::now() - started_ms) >= timeout_ms as f64;
if timed_out {
return Err(anyhow::anyhow!(
"peer {} has no active transport for protocol traffic",
id
));
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(50)).await;
}
}
async fn resolve_diagnostic_peer_endpoint(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(PeerSessionSnapshot, iroh::EndpointId)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self
.peer_session(id)
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
let connection_id = snapshot
.active_connection_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("peer {} is missing an active connection", id))?;
if let Some((_rejected, reason)) = self.session_admission_block_reason(connection_id) {
return Err(anyhow::anyhow!(
"peer {} is not admitted for diagnostic traffic: {}",
id,
reason
));
}
if self.connection_requires_application_crypto(connection_id)
&& self
.application_crypto_key_for_connection(Some(connection_id))
.is_none()
{
return Err(anyhow::anyhow!(
"peer {} requires application crypto but no key is installed",
id
));
}
let node_id = snapshot
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("peer {} is missing a remote node id", id))?;
let endpoint_id = node_id.parse::<iroh::EndpointId>().map_err(|error| {
anyhow::anyhow!("peer {} has invalid remote node id: {}", id, error)
})?;
if self.is_connection_transport_alive(endpoint_id).await {
return Ok((snapshot, endpoint_id));
}
let _ = self
.confirm_managed_connection_readiness(connection_id)
.await;
#[cfg(not(target_arch = "wasm32"))]
let timed_out = started.elapsed().as_millis() >= timeout_ms as u128;
#[cfg(target_arch = "wasm32")]
let timed_out = (js_sys::Date::now() - started_ms) >= timeout_ms as f64;
if timed_out {
return Err(anyhow::anyhow!(
"peer {} has no active transport for diagnostic traffic",
id
));
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(50)).await;
}
}
pub async fn add_peer_scope(&self, id: &str, scope: &str) -> Vec<String> {
self.connection_manager.add_scope(id, scope).await
}
pub async fn release_peer_scope(&self, id: &str, scope: Option<&str>) -> Vec<String> {
self.connection_manager.release_scope(id, scope).await
}
pub async fn peer_scopes(&self, id: &str) -> Vec<String> {
self.connection_manager.get_scopes(id).await
}
pub async fn same_peer(&self, left: &str, right: &str) -> bool {
self.connection_manager.are_same_peer(left, right).await
}
pub async fn peer_snapshot(&self, id: &str) -> Option<crate::connection_manager::PeerSnapshot> {
if let Some(snapshot) = self.refresh_cached_peer_snapshot(id).await {
return Some(snapshot);
}
let _ = self
.adopt_live_transport_session(id, "peer-snapshot-query")
.await;
self.connection_manager.peer_snapshot(id).await
}
pub async fn peer_session(&self, id: &str) -> Option<PeerSessionSnapshot> {
if let Some(snapshot) = self
.refresh_cached_peer_snapshot(id)
.await
.map(peer_session_snapshot_from_peer)
{
return Some(self.apply_peer_session_admission_guard(snapshot));
}
self.adopt_live_transport_session(id, "peer-session-query")
.await
.map(|snapshot| self.apply_peer_session_admission_guard(snapshot))
}
pub async fn peer_sessions(&self) -> Vec<PeerSessionSnapshot> {
let mut refreshed_snapshots = Vec::new();
for snapshot in self.connection_manager.list_peer_snapshots().await {
let lookup_id = snapshot.peer_id.clone();
if let Some(refreshed) = self.refresh_cached_peer_snapshot(&lookup_id).await {
refreshed_snapshots.push(refreshed);
}
}
let snapshots: Vec<PeerSessionSnapshot> = refreshed_snapshots
.into_iter()
.map(peer_session_snapshot_from_peer)
.map(|snapshot| self.apply_peer_session_admission_guard(snapshot))
.collect();
#[cfg(target_arch = "wasm32")]
if snapshots.is_empty() {
let active_endpoints = {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.active_endpoint_ids().await.len()
} else {
0
}
};
if active_endpoints > 0 {
let now_ms = js_sys::Date::now() as u64;
let previous_ms = self
.last_empty_peer_sessions_warning_ms
.load(std::sync::atomic::Ordering::Relaxed);
if now_ms.saturating_sub(previous_ms)
>= Self::EMPTY_PEER_SESSIONS_WARNING_INTERVAL_MS
{
self.last_empty_peer_sessions_warning_ms
.store(now_ms, std::sync::atomic::Ordering::Relaxed);
let record_count = self.connection_manager.list_all().await.len();
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][peer-sessions][wasm] empty snapshot list active_endpoints={} record_count={} note=expected briefly during token-first connects before the main stream is registered; live-transports-may-belong-to-another-client-instance-until-adopted",
active_endpoints, record_count
)));
}
}
}
snapshots
}
pub async fn connection_state(&self, connection_id: &str) -> Option<ConnectionStateSnapshot> {
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await?;
#[allow(unused_mut)]
let mut peer_snapshot = self.connection_manager.peer_snapshot(connection_id).await;
Some(
self.apply_connection_state_admission_guard(connection_state_snapshot_from_parts(
&record,
peer_snapshot.as_ref(),
)),
)
}
pub async fn connection_states(&self) -> Vec<ConnectionStateSnapshot> {
let records = self.connection_manager.list_active().await;
let mut snapshots = Vec::with_capacity(records.len());
for record in records {
let peer_snapshot = self
.connection_manager
.peer_snapshot(&record.connection_id)
.await;
snapshots.push(self.apply_connection_state_admission_guard(
connection_state_snapshot_from_parts(&record, peer_snapshot.as_ref()),
));
}
snapshots
}
#[cfg(target_arch = "wasm32")]
pub async fn emit_current_wasm_connection_state(&self, connection_id: &str) {
if let Some(snapshot) = self.connection_state(connection_id).await {
let fingerprint = super::WasmConnectionStateFingerprint::from(&snapshot);
let should_emit = {
let mut emitted = self
.last_emitted_connection_states
.lock()
.expect("wasm connection state cache poisoned");
match emitted.get(connection_id) {
Some(previous) if *previous == fingerprint => false,
_ => {
emitted.insert(connection_id.to_string(), fingerprint);
true
}
}
};
if !should_emit {
return;
}
emit_wasm_connection_state_event(&snapshot);
}
}
pub async fn wait_for_settled_peer(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> Option<PeerSessionSnapshot> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self.peer_session(id).await;
if let Some(snapshot_ref) = snapshot.as_ref() {
if snapshot_ref.settled_ready {
let stable_deadline_ms = snapshot_ref
.last_lifecycle_transition_at_ms
.saturating_add(Self::SETTLED_PEER_STABILITY_WINDOW_MS as i64);
#[cfg(not(target_arch = "wasm32"))]
let stable_now = now_millis_i64() >= stable_deadline_ms;
#[cfg(target_arch = "wasm32")]
let stable_now = js_sys::Date::now() as i64 >= stable_deadline_ms;
if stable_now {
return snapshot;
}
} else if let Some(connection_id) = snapshot_ref.active_connection_id.clone() {
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
}
}
#[cfg(not(target_arch = "wasm32"))]
let timed_out = started.elapsed().as_millis() >= timeout_ms as u128;
#[cfg(target_arch = "wasm32")]
let timed_out = (js_sys::Date::now() - started_ms) >= timeout_ms as f64;
if timed_out {
return snapshot;
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(50)).await;
}
}
pub(crate) async fn confirm_managed_connection_readiness(&self, connection_id: &str) -> bool {
let healthy = self.probe_peer_health(connection_id).await;
#[cfg(target_arch = "wasm32")]
self.emit_current_wasm_connection_state(connection_id).await;
healthy
}
pub async fn open_peer_bi(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
let (snapshot, endpoint_id) = self.resolve_settled_peer_endpoint(id, timeout_ms).await?;
let remote_node_id = endpoint_id.to_string();
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
let node_guard = self.iroh_node.read().await;
let node = node_guard
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let (send, recv) = node.open_bi(endpoint_id).await?;
let key = self
.application_crypto_key_for_connection_or_endpoint(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
)
.await;
crate::application_crypto_streams::wrap_peer_streams(key, send, recv)
.map_err(|error| anyhow::anyhow!("application crypto stream wrap failed: {error}"))
.map(|wrapped| {
(
snapshot.active_connection_id,
remote_node_id,
wrapped.0,
wrapped.1,
)
})
}
pub async fn open_peer_bi_explicit_file_sender(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
)> {
let (snapshot, endpoint_id) = self.resolve_settled_peer_endpoint(id, timeout_ms).await?;
let remote_node_id = endpoint_id.to_string();
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
let node_guard = self.iroh_node.read().await;
let node = node_guard
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let (mut send, _recv) = node.open_bi(endpoint_id).await?;
tokio::io::AsyncWriteExt::write_all(
&mut send,
&[crate::explicit_transfer_crypto::EXPLICIT_FILE_PROTOCOL_BYTE],
)
.await?;
let connection_id = snapshot.active_connection_id.clone();
let key = self
.application_crypto_key_for_connection_or_endpoint(
connection_id.as_deref(),
&endpoint_id,
)
.await;
Ok((
connection_id.clone(),
remote_node_id,
crate::application_crypto_streams::wrap_peer_send_stream(key, send),
))
}
pub async fn open_peer_bi_diagnostic(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
let (snapshot, endpoint_id) = self
.resolve_diagnostic_peer_endpoint(id, timeout_ms)
.await?;
let remote_node_id = endpoint_id.to_string();
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
let node_guard = self.iroh_node.read().await;
let node = node_guard
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let (send, recv) = node.open_bi(endpoint_id).await?;
self.wrap_peer_streams_for_connection(snapshot.active_connection_id.as_deref(), send, recv)
.map(|wrapped| {
(
snapshot.active_connection_id,
remote_node_id,
wrapped.0,
wrapped.1,
)
})
}
pub async fn open_peer_bi_transport_only(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
iroh::endpoint::SendStream,
iroh::endpoint::RecvStream,
)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self
.peer_session(id)
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
let node_id = snapshot
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("peer {} is missing a remote node id", id))?;
let endpoint_id = node_id.parse::<iroh::EndpointId>().map_err(|error| {
anyhow::anyhow!("peer {} has invalid remote node id: {}", id, error)
})?;
if self.is_connection_transport_alive(endpoint_id).await {
if self
.application_crypto_key_for_endpoint(&endpoint_id)
.await
.is_some()
{
anyhow::bail!(
"open_peer_bi_transport_only is not allowed while application crypto is active; use open_peer_bi instead"
);
}
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
let node_guard = self.iroh_node.read().await;
let node = node_guard
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let (send, recv) = node.open_bi(endpoint_id).await?;
return Ok((
snapshot.active_connection_id,
endpoint_id.to_string(),
send,
recv,
));
}
#[cfg(not(target_arch = "wasm32"))]
let timed_out = started.elapsed().as_millis() >= timeout_ms as u128;
#[cfg(target_arch = "wasm32")]
let timed_out = (js_sys::Date::now() - started_ms) >= timeout_ms as f64;
if timed_out {
return Err(anyhow::anyhow!(
"peer {} transport not alive (transport-only open timed out)",
id
));
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(50)).await;
}
}
pub async fn open_peer_uni(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
)> {
let (snapshot, endpoint_id) = self.resolve_settled_peer_endpoint(id, timeout_ms).await?;
let remote_node_id = endpoint_id.to_string();
let node_guard = self.iroh_node.read().await;
let node = node_guard
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let send = node.open_uni(endpoint_id).await?;
let connection_id = snapshot.active_connection_id.clone();
Ok((
connection_id,
remote_node_id,
self.wrap_peer_send_stream_for_connection(
snapshot.active_connection_id.as_deref(),
send,
),
))
}
pub async fn resolve_peer_connection_ids(&self, id: &str) -> Vec<String> {
let mut connection_ids = Vec::new();
if let Some(snapshot) = self.connection_manager.peer_snapshot(id).await {
connection_ids.extend(snapshot.connection_ids);
}
if let Some(record) = self.connection_manager.get_by_connection_id(id).await {
connection_ids.push(record.connection_id);
}
if let Some((left, right)) = id.split_once('-') {
for node_id in [left.trim(), right.trim()] {
if node_id.is_empty() {
continue;
}
for record in self.connection_manager.get_by_node_id(node_id).await {
if let Some(snapshot) = self
.connection_manager
.peer_snapshot(&record.connection_id)
.await
{
connection_ids.extend(snapshot.connection_ids);
} else {
connection_ids.push(record.connection_id);
}
}
}
let all_records = self.connection_manager.list_all().await;
for node_id in [left.trim(), right.trim()] {
if node_id.is_empty() {
continue;
}
for record in all_records.iter().filter(|record| {
record
.connection_id
.split('-')
.any(|part| part.trim() == node_id)
}) {
if let Some(snapshot) = self
.connection_manager
.peer_snapshot(&record.connection_id)
.await
{
connection_ids.extend(snapshot.connection_ids);
} else {
connection_ids.push(record.connection_id.clone());
}
}
}
}
if connection_ids.len() > 1 {
let mut unique = std::collections::HashSet::with_capacity(connection_ids.len());
connection_ids.retain(|connection_id| unique.insert(connection_id.clone()));
}
connection_ids
}
pub async fn resolve_peer_connection_records(
&self,
id: &str,
) -> Vec<crate::connection_manager::ConnectionRecord> {
let mut records = Vec::new();
for connection_id in self.resolve_peer_connection_ids(id).await {
if let Some(record) = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
records.push(record);
}
}
records.sort_by(|left, right| right.updated_at_ms.cmp(&left.updated_at_ms));
records
}
pub async fn best_connection_record_for_peer(
&self,
id: &str,
) -> Option<crate::connection_manager::ConnectionRecord> {
self.connection_manager.best_connection_for_peer(id).await
}
pub async fn list_managed_connections(
&self,
) -> Vec<crate::connection_manager::ConnectionRecord> {
self.connection_manager.list_all().await
}
pub async fn managed_connection_device_hint(&self, id: &str) -> Option<String> {
self.resolve_peer_connection_records(id)
.await
.into_iter()
.find_map(|record| {
record
.device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
record
.device_id_hint
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
})
})
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn managed_connection_adoption(
&self,
record: &crate::connection_manager::ConnectionRecord,
) -> Option<crate::client::ManagedConnectionAdoption> {
if !matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) {
return None;
}
let admitted_scope = if self.session_registry_active() {
let admission = self.session_token_registry.admission(&record.connection_id);
match &admission {
crate::session_token::SessionAdmission::Accepted { scope, .. } => scope.clone(),
crate::session_token::SessionAdmission::Pending => return None,
crate::session_token::SessionAdmission::Rejected { reason } => {
eprintln!(
"[PlutoRTC][managed-adoption][blocked] connection_id={} reason={} node_id={} device_id_hint={}",
record.connection_id,
reason,
record.node_id.as_deref().unwrap_or("pending"),
record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.unwrap_or("pending")
);
return None;
}
}
} else {
None
};
let node_id = record
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())?
.to_string();
let device_id = self
.managed_connection_device_hint(&record.connection_id)
.await;
let resolved_device_id = device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
if matches!(
admitted_scope.as_ref().map(|scope| scope.as_str()),
Some("user-device")
) {
record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
} else {
None
}
})
.or_else(|| {
if matches!(
admitted_scope.as_ref().map(|scope| scope.as_str()),
Some("share") | Some("magic-link")
) {
Some(node_id.clone())
} else {
None
}
});
if resolved_device_id.is_none()
&& matches!(
admitted_scope.as_ref().map(|scope| scope.as_str()),
Some("user-device")
)
{
println!(
"[PlutoRTC][managed-adoption][wait] connection_id={} scope=user-device reason=missing-authoritative-device-id node_id={} device_id_hint={}",
record.connection_id,
node_id,
record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.unwrap_or("pending")
);
}
Some(crate::client::ManagedConnectionAdoption {
connection_id: record.connection_id.clone(),
node_id,
device_id: resolved_device_id.clone(),
transport_generation: record.transport_generation,
status_reason: record.status_reason.clone(),
main_stream_ready: resolved_device_id.is_some(),
})
}
pub async fn retire_managed_connection(&self, connection_id: &str, reason: Option<String>) {
let normalized_reason = reason.and_then(|value| {
let trimmed = value.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
});
#[cfg(target_arch = "wasm32")]
{
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][managed-retire] connection_id={} reason={:?}",
connection_id, normalized_reason
)));
}
#[cfg(not(target_arch = "wasm32"))]
{
println!(
"[pluto-rtc][managed-retire] connection_id={} reason={:?}",
connection_id, normalized_reason
);
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
if self
.defer_managed_retirement_while_webrtc_active(
connection_id,
normalized_reason.clone(),
)
.await
{
return;
}
}
self.retire_managed_connection_now(connection_id, normalized_reason)
.await;
}
pub async fn bind_connection_device_id(
&self,
connection_id: &str,
device_id: &str,
) -> Option<crate::connection_manager::PeerSnapshot> {
let trimmed = device_id.trim();
if !trimmed.is_empty() {
self.mark_trusted_user_device_connection_admitted(connection_id, trimmed)
.await;
if let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
{
if let Some(node_id) = record.node_id.as_deref() {
self.prune_stale_records_for_device_node(trimmed, node_id)
.await;
}
}
}
self.connection_manager.peer_snapshot(connection_id).await
}
pub async fn report_managed_connection_settled(
&self,
connection_id: &str,
settled: bool,
device_id: Option<&str>,
) -> Option<crate::connection_manager::PeerSnapshot> {
if let Some(device_id) = device_id.map(str::trim).filter(|value| !value.is_empty()) {
let _ = self
.connection_manager
.set_device_id(connection_id, device_id.to_string())
.await;
}
self.connection_manager
.set_health(
connection_id,
if settled {
crate::connection_manager::ConnectionHealth::Healthy
} else {
crate::connection_manager::ConnectionHealth::Unknown
},
)
.await
}
pub async fn probe_peer_health(&self, id: &str) -> bool {
let snapshot = self.connection_manager.peer_snapshot(id).await;
let Some(snapshot) = snapshot else {
return false;
};
let Some(node_id) = snapshot.node_id.clone() else {
let _ = self
.connection_manager
.set_health(id, crate::connection_manager::ConnectionHealth::Stale)
.await;
return false;
};
let Ok(endpoint_id) = node_id.parse::<iroh::EndpointId>() else {
let _ = self
.connection_manager
.set_health(id, crate::connection_manager::ConnectionHealth::Stale)
.await;
return false;
};
#[cfg(not(target_arch = "wasm32"))]
let healthy = self.is_connection_transport_alive(endpoint_id).await;
#[cfg(target_arch = "wasm32")]
let healthy = self.is_connection_transport_alive(endpoint_id).await;
let health = if healthy {
crate::connection_manager::ConnectionHealth::Healthy
} else {
crate::connection_manager::ConnectionHealth::Stale
};
let _ = self.connection_manager.set_health(id, health).await;
healthy
}
pub async fn update_presence(
&self,
user_id: &str,
device_name: &str,
ticket: &str,
metadata: Option<&str>,
) -> anyhow::Result<()> {
self.update_presence_with_ttl(user_id, device_name, ticket, 300_000, metadata)
.await
}
pub async fn update_presence_with_ttl(
&self,
user_id: &str,
device_name: &str,
ticket: &str,
ttl_ms: u64,
metadata: Option<&str>,
) -> anyhow::Result<()> {
let effective_ticket = self.refresh_presence_ticket_for_publication(ticket).await;
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(&effective_ticket);
parse_endpoint_ticket(iroh_ticket)?;
let node_guard = self.node_id.read().await;
#[cfg(not(target_arch = "wasm32"))]
let metadata_owned = self.metadata_with_local_device_id(metadata).await;
#[cfg(not(target_arch = "wasm32"))]
let metadata = metadata_owned.as_deref();
if let Some(node_id) = node_guard.as_deref() {
self.signaling
.update_presence(
user_id,
node_id,
&effective_ticket,
true,
device_name,
ttl_ms,
metadata,
)
.await?;
return Ok(());
}
#[cfg(not(target_arch = "wasm32"))]
{
let (_iroh_fingerprint, scope, token_fingerprint) =
super::core_impl::summarize_compound_ticket_for_logs(effective_ticket.as_str());
eprintln!(
"[pluto-rtc][presence][publish][blocked] user_id={} scope={} token_fp={} reason=node-id-uninitialized",
user_id,
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string())
);
}
Err(anyhow::anyhow!(
"Iroh node not initialized for native presence publish"
))
}
async fn refresh_presence_ticket_for_publication(&self, ticket: &str) -> String {
let trimmed = ticket.trim();
if trimmed.is_empty() {
return ticket.to_string();
}
let (iroh_ticket, suffix) = crate::session_token::split_compound_ticket(trimmed);
let Some(payload) = suffix.and_then(|value| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, value)
}) else {
return trimmed.to_string();
};
let latest_iroh_ticket = match self.endpoint_ticket().await {
Ok(value) => value,
Err(_) => return trimmed.to_string(),
};
if latest_iroh_ticket == iroh_ticket {
return trimmed.to_string();
}
let rebuilt = crate::session_token::build_compound_ticket(
&latest_iroh_ticket,
&payload.token,
payload.scope.clone(),
payload.max_connections,
);
#[cfg(not(target_arch = "wasm32"))]
{
let (iroh_fingerprint, scope, token_fingerprint) =
super::core_impl::summarize_compound_ticket_for_logs(rebuilt.as_str());
println!(
"[pluto-rtc][presence][ticket-refresh] rebuilt compound ticket for publish scope={} endpoint_changed=true token_fp={} iroh_fp={}",
scope.unwrap_or_else(|| payload.scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint
);
}
#[cfg(target_arch = "wasm32")]
web_sys::console::info_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][presence] rebuilt compound ticket for publish scope={} endpoint_changed=true",
payload.scope
)));
rebuilt
}
pub async fn set_offline(&self, user_id: &str) -> anyhow::Result<()> {
let node_guard = self.node_id.read().await;
if let Some(node_id) = node_guard.as_deref() {
self.signaling.set_offline(user_id, node_id).await?;
}
Ok(())
}
pub async fn update_device(
&self,
user_id: &str,
device_id: &str,
device_name: Option<&str>,
capabilities: Option<crate::signaling::DeviceCapabilities>,
metadata: Option<&str>,
) -> anyhow::Result<()> {
self.signaling
.update_device(user_id, device_id, device_name, capabilities, metadata)
.await
}
pub async fn delete_device(&self, user_id: &str, device_id: &str) -> anyhow::Result<()> {
self.signaling.delete_device(user_id, device_id).await
}
pub fn active_session_identity(&self) -> Option<(String, String)> {
match self.auto_connect_loop_key.lock() {
Ok(guard) => guard.clone(),
Err(p) => p.into_inner().clone(),
}
}
pub async fn update_signaling_excluded_peers(
&self,
user_id: &str,
excluded_peers: &[String],
) -> anyhow::Result<()> {
let node_guard = self.node_id.read().await;
if let Some(node_id) = node_guard.as_deref() {
self.signaling
.set_excluded_peers(user_id, node_id, excluded_peers)
.await?;
}
Ok(())
}
pub fn current_excluded_peers_snapshot(&self) -> Vec<String> {
let guard = match self.auto_connect_excluded.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
let mut peers: Vec<String> = guard.iter().cloned().collect();
peers.sort();
peers
}
fn normalize_auto_connect_exclusion_key(value: &str) -> Option<String> {
let value = value.trim().to_ascii_lowercase();
if value.is_empty() {
None
} else {
Some(value)
}
}
pub async fn exclude_peer_and_publish(self: &std::sync::Arc<Self>, remote_device_id: &str) {
self.set_auto_connect_excluded_peer(remote_device_id, None, true);
let Some((user_id, _)) = self.active_session_identity() else {
return;
};
let excluded = self.current_excluded_peers_snapshot();
if let Err(error) = self
.update_signaling_excluded_peers(&user_id, &excluded)
.await
{
eprintln!(
"[PlutoRTC] exclude_peer_and_publish failed to publish excluded_peers remote_device_id={} error={}",
remote_device_id,
error
);
}
}
pub async fn unexclude_peer_and_publish(self: &std::sync::Arc<Self>, remote_device_id: &str) {
self.set_auto_connect_excluded_peer(remote_device_id, None, false);
let Some((user_id, _)) = self.active_session_identity() else {
return;
};
let excluded = self.current_excluded_peers_snapshot();
if let Err(error) = self
.update_signaling_excluded_peers(&user_id, &excluded)
.await
{
eprintln!(
"[PlutoRTC] unexclude_peer_and_publish failed to publish excluded_peers remote_device_id={} error={}",
remote_device_id,
error
);
}
}
pub async fn search_devices(
&self,
user_id: &str,
) -> anyhow::Result<Vec<crate::signaling::Device>> {
let node_guard = self.node_id.read().await;
self.signaling
.search_devices(user_id, node_guard.as_deref())
.await
}
pub async fn devices_with_status(
&self,
user_id: &str,
) -> anyhow::Result<Vec<DeviceStatusSnapshot>> {
let node_guard = self.node_id.read().await;
let devices = self
.signaling
.list_devices(user_id, node_guard.as_deref())
.await?;
self.refresh_device_status_peer_snapshots_from_devices(&devices)
.await;
let peers = self.connection_manager.list_peer_snapshots().await;
self.backfill_peer_device_ids_from_devices(&devices, &peers)
.await;
let peers = self.connection_manager.list_peer_snapshots().await;
let _discovered_count = devices.len();
let _peer_count = peers.len();
let merged = merge_device_status_snapshots(devices, peers);
Ok(merged
.into_iter()
.map(|snapshot| self.apply_device_status_admission_guard(snapshot))
.collect())
}
pub(super) async fn refresh_device_status_peer_snapshots_from_devices(
&self,
devices: &[crate::signaling::Device],
) {
let mut seen = std::collections::HashSet::new();
for device in devices {
let lookup_ids = [
device
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned),
Some(device.device_id.trim().to_string()).filter(|value| !value.is_empty()),
];
for lookup_id in lookup_ids.into_iter().flatten() {
if !seen.insert(lookup_id.clone()) {
continue;
}
if self
.refresh_cached_peer_snapshot(&lookup_id)
.await
.is_some()
{
continue;
}
let _ = self
.adopt_live_transport_session(&lookup_id, "devices-with-status")
.await;
}
}
}
pub(super) async fn backfill_peer_device_ids_from_devices(
&self,
devices: &[crate::signaling::Device],
peers: &[crate::connection_manager::PeerSnapshot],
) {
let device_ids_by_node: std::collections::HashMap<String, String> = devices
.iter()
.filter_map(|device| {
let node_id = device
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())?;
let device_id = device.device_id.trim();
if device_id.is_empty() {
return None;
}
Some((node_id.to_string(), device_id.to_string()))
})
.collect();
for peer in peers {
let has_device_id = peer
.device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.is_some();
if has_device_id {
continue;
}
let Some(node_id) = peer
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
else {
continue;
};
let Some(device_id) = device_ids_by_node.get(node_id).cloned() else {
continue;
};
for connection_id in &peer.connection_ids {
let current_device_id = self
.connection_manager
.get_by_connection_id(connection_id)
.await
.and_then(|record| {
record
.device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
});
if current_device_id.as_deref() == Some(device_id.as_str()) {
continue;
}
self.connection_manager
.set_device_id(connection_id, device_id.clone())
.await;
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][devices-with-status] backfilled peer device_id connection_id={} node_id={} device_id={}",
connection_id, node_id, device_id
)));
}
}
}
pub async fn runtime_status(&self) -> RuntimeStatus {
let node_id = self.node_id.read().await.clone();
let transport_config = self.transport_config.read().await;
let transport = RuntimeTransportStatus::from_config(&transport_config);
RuntimeStatus {
ready: true,
node_id,
transport,
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn notify_network_change(&self) -> anyhow::Result<usize> {
let endpoint = {
let endpoint_guard = self.iroh_endpoint.read().await;
endpoint_guard
.clone()
.ok_or_else(|| anyhow::anyhow!("Iroh endpoint not initialized"))?
};
tokio::time::timeout(
std::time::Duration::from_secs(10),
endpoint.network_change(),
)
.await
.map_err(|_| anyhow::anyhow!("Iroh endpoint network-change refresh timed out after 10s"))?;
let records = self.connection_manager.list_all().await;
let mut retired = 0usize;
for record in records {
let Some(endpoint_id) = record
.endpoint_id
.as_deref()
.or(record.node_id.as_deref())
.and_then(|value| value.parse::<iroh::EndpointId>().ok())
else {
continue;
};
let transport_alive = tokio::time::timeout(
std::time::Duration::from_secs(5),
self.is_connection_transport_alive(endpoint_id),
)
.await
.map_err(|_| {
anyhow::anyhow!(
"Iroh transport health check timed out after 5s endpoint_id={endpoint_id}"
)
})?;
if transport_alive {
continue;
}
self.connection_manager
.set_closing(
&record.connection_id,
Some("network-change-retired-stale-record".to_string()),
)
.await;
self.connection_manager
.set_closed(
&record.connection_id,
Some("network-change-retired-stale-record".to_string()),
)
.await;
let _ = self.connection_manager.remove(&record.connection_id).await;
retired = retired.saturating_add(1);
}
Ok(retired)
}
pub async fn search_devices_raw(
&self,
user_id: &str,
) -> anyhow::Result<Vec<crate::signaling::Device>> {
self.signaling.search_devices(user_id, None).await
}
pub async fn connect_device(
&self,
device_id: Option<&str>,
endpoint_ticket: &str,
) -> anyhow::Result<ManagedConnectResult> {
let (iroh_ticket, token_suffix) =
crate::session_token::split_compound_ticket(endpoint_ticket.trim());
let extracted_token_payload = token_suffix.and_then(|suffix| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, suffix)
});
let extracted_token = extracted_token_payload
.as_ref()
.map(|payload| payload.token.clone());
let extracted_token_suffix = token_suffix
.filter(|_| extracted_token_payload.is_some())
.map(ToOwned::to_owned);
if token_suffix.is_some() && extracted_token_payload.is_none() {
return Err(anyhow::anyhow!(
"invalid compound ticket payload for endpoint ticket"
));
}
let endpoint_addr = parse_endpoint_ticket(iroh_ticket)?;
let remote_node_id = endpoint_addr.id.to_string();
let local_node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
let trimmed_device_id = device_id
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned);
let connect_gate = self.managed_connect_gate(&connection_id).await;
let _connect_guard = connect_gate.lock().await;
let mut approved_scope: Option<String> = None;
let log_phase = |phase: &str| {
let message = format!(
"[pluto-rtc][connect-device][phase] phase={} connection_id={} remote_node_id={} device_id={}",
phase,
connection_id,
remote_node_id,
trimmed_device_id.as_deref().unwrap_or("")
);
#[cfg(not(target_arch = "wasm32"))]
println!("{}", message);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&message));
};
log_phase("gate-acquired");
if let Some(existing) = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
if matches!(
existing.state,
crate::connection_manager::ConnectionState::Connected
) && self.is_connection_transport_alive(endpoint_addr.id).await
{
if let Some(device_id) = trimmed_device_id.clone() {
self.mark_trusted_user_device_connection_admitted(&connection_id, &device_id)
.await;
self.prune_stale_records_for_device_node(&device_id, &remote_node_id)
.await;
}
let needs_token_representation =
should_represent_session_token_on_connected_transport(
extracted_token.is_some(),
self.session_admission_block_reason(&connection_id)
.is_some(),
);
if needs_token_representation {
if let Some(token) = extracted_token.as_ref() {
println!(
"[PlutoRTC] Existing managed transport still awaiting session admission; re-presenting token connection_id={} remote_node_id={} token_fp={}",
connection_id,
remote_node_id,
super::core_impl::log_fingerprint(token.as_str())
);
match self
.present_session_token_to_host_with_payload(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
)
.await
{
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
self.mark_connection_admitted_by_host(
&connection_id,
Some(approval_scope.as_str()),
trimmed_device_id.clone(),
)
.await;
println!(
"[PlutoRTC] Re-presented session-token on existing managed transport connection_id={} scope={} token_fp={} local_admission=accepted",
connection_id,
approval_scope,
super::core_impl::log_fingerprint(token.as_str())
);
}
Err(error) => {
eprintln!(
"[PlutoRTC] Failed to re-present session-token on existing managed transport connection_id={} error={}",
connection_id,
error
);
}
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
}
}
let snapshot = self.connection_manager.peer_snapshot(&connection_id).await;
return Ok(ManagedConnectResult {
connection_id,
device_id: snapshot
.as_ref()
.and_then(|value| value.device_id.clone())
.or(trimmed_device_id.clone()),
device_id_hint: snapshot
.as_ref()
.and_then(|value| value.device_id_hint.clone())
.or(trimmed_device_id),
remote_node_id,
state: "connected".to_string(),
approved_scope,
});
}
}
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
trimmed_device_id.clone(),
Some(remote_node_id.clone()),
)
.await;
if let Some(device_id) = trimmed_device_id.as_deref() {
self.prune_stale_records_for_device_node(device_id, &remote_node_id)
.await;
}
self.connection_manager.set_connecting(&connection_id).await;
log_phase("transport-dial-start");
self.ensure_connected_addr(endpoint_addr.id, endpoint_addr.clone())
.await?;
log_phase("transport-dial-complete");
let transport_stable_id = self
.get_connection(endpoint_addr.id)
.await
.map(|connection| connection.stable_id() as u64);
log_phase("transport-record-start");
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.clone()),
transport_stable_id,
None,
)
.await;
log_phase("transport-record-complete");
log_phase("transport-health-start");
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
log_phase("transport-health-complete");
if let Some(device_id) = trimmed_device_id.clone() {
log_phase("trusted-admission-start");
self.mark_trusted_user_device_connection_admitted(&connection_id, &device_id)
.await;
log_phase("trusted-admission-complete");
log_phase("stale-record-prune-start");
self.prune_stale_records_for_device_node(&device_id, &remote_node_id)
.await;
log_phase("stale-record-prune-complete");
}
if let Some(ref token) = extracted_token {
match self
.present_session_token_to_host_with_payload(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
)
.await
{
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
self.mark_connection_admitted_by_host(
&connection_id,
Some(approval_scope.as_str()),
trimmed_device_id.clone(),
)
.await;
println!(
"[PlutoRTC] Sent session-token presentation to host and received approval connection_id={} scope={} token_fp={} local_admission=accepted",
connection_id,
approval_scope,
super::core_impl::log_fingerprint(token.as_str())
);
}
Err(error) => {
eprintln!(
"[PlutoRTC] Failed to send session-token presentation connection_id={} token_fp={} error={}",
connection_id,
super::core_impl::log_fingerprint(token.as_str()),
error
);
}
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
}
let snapshot = self.connection_manager.peer_snapshot(&connection_id).await;
let state = match snapshot.as_ref().map(|value| &value.status) {
Some(crate::connection_manager::ConnectionState::Pending)
| Some(crate::connection_manager::ConnectionState::Connecting) => "connecting",
Some(crate::connection_manager::ConnectionState::Connected) => "connected",
Some(crate::connection_manager::ConnectionState::Failed) => "failed",
Some(crate::connection_manager::ConnectionState::Closed)
| Some(crate::connection_manager::ConnectionState::Closing) => "closed",
None => "connected",
}
.to_string();
log_phase("return");
Ok(ManagedConnectResult {
connection_id,
device_id: snapshot
.as_ref()
.and_then(|value| value.device_id.clone())
.or(trimmed_device_id.clone()),
device_id_hint: snapshot
.as_ref()
.and_then(|value| value.device_id_hint.clone())
.or(trimmed_device_id),
remote_node_id,
state,
approved_scope,
})
}
pub async fn managed_connection_health(
&self,
connection_id: &str,
) -> Option<ManagedConnectionHealthSnapshot> {
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await?;
#[allow(unused_mut)]
let mut peer_snapshot = self.connection_manager.peer_snapshot(connection_id).await;
#[cfg(not(target_arch = "wasm32"))]
let transport_alive = matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) && self
.is_managed_connection_transport_alive(connection_id)
.await;
#[cfg(not(target_arch = "wasm32"))]
if transport_alive
&& !peer_snapshot
.as_ref()
.map(peer_snapshot_settled_ready)
.unwrap_or(false)
{
let refreshed = self
.confirm_managed_connection_readiness(connection_id)
.await;
if refreshed {
peer_snapshot = self.connection_manager.peer_snapshot(connection_id).await;
}
}
let settled_ready = peer_snapshot
.as_ref()
.map(peer_snapshot_settled_ready)
.unwrap_or(false);
let admission_block = self.session_admission_block_reason(connection_id);
let admitted_settled_ready = settled_ready && admission_block.is_none();
let settle_deadline_elapsed = matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) && !admitted_settled_ready
&& now_millis_i64().saturating_sub(record.last_transport_change_at_ms)
>= MANAGED_SETTLE_DEADLINE_MS;
let status = match record.state {
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting => {
ManagedConnectionHealthStatus::AwaitingReplacement
}
crate::connection_manager::ConnectionState::Connected => {
#[cfg(not(target_arch = "wasm32"))]
{
if transport_alive && admitted_settled_ready {
ManagedConnectionHealthStatus::Healthy
} else if settle_deadline_elapsed {
ManagedConnectionHealthStatus::Dead
} else {
ManagedConnectionHealthStatus::AwaitingReplacement
}
}
#[cfg(target_arch = "wasm32")]
{
let peer_health = peer_snapshot
.as_ref()
.map(|snapshot| snapshot.health.clone());
match peer_health
.unwrap_or(crate::connection_manager::ConnectionHealth::Unknown)
{
crate::connection_manager::ConnectionHealth::Healthy
if admitted_settled_ready =>
{
ManagedConnectionHealthStatus::Healthy
}
_ if settle_deadline_elapsed => ManagedConnectionHealthStatus::Dead,
crate::connection_manager::ConnectionHealth::Unknown
| crate::connection_manager::ConnectionHealth::Suspect
| crate::connection_manager::ConnectionHealth::Stale
| crate::connection_manager::ConnectionHealth::Healthy => {
ManagedConnectionHealthStatus::AwaitingReplacement
}
}
}
}
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
| crate::connection_manager::ConnectionState::Failed => {
ManagedConnectionHealthStatus::Dead
}
};
let readiness_state = connection_readiness_state(&record, peer_snapshot.as_ref());
let readiness_reason = connection_readiness_reason(&record, peer_snapshot.as_ref());
let replacement_pending = matches!(readiness_state, ReadinessState::AwaitingReplacement);
let last_lifecycle_transition_at_ms =
record.last_transport_change_at_ms.max(record.updated_at_ms);
let mut snapshot = ManagedConnectionHealthSnapshot {
connection_id: record.connection_id,
device_id: record.device_id,
device_id_hint: record.device_id_hint,
node_id: record.node_id,
active_transport_stable_id: record.transport_stable_id,
transport_generation: record.transport_generation,
status,
settled_ready: admitted_settled_ready,
readiness_state,
replacement_pending,
last_lifecycle_transition_at_ms,
readiness_reason,
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,
last_reconnect_reason: record.last_reconnect_reason,
};
if let Some((rejected, reason)) = admission_block {
snapshot.settled_ready = false;
snapshot.replacement_pending = false;
snapshot.readiness_reason = reason.clone();
if rejected {
snapshot.status = ManagedConnectionHealthStatus::Dead;
snapshot.readiness_state = ReadinessState::Failed;
} else {
snapshot.status = ManagedConnectionHealthStatus::AwaitingReplacement;
snapshot.readiness_state = ReadinessState::Settling;
}
}
Some(snapshot)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn managed_connection_bridge_action(
&self,
connection_id: &str,
) -> Option<ManagedConnectionBridgeAction> {
let snapshot = self.managed_connection_health(connection_id).await?;
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await?;
let settle_deadline_elapsed = matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) && !snapshot.settled_ready
&& now_millis_i64().saturating_sub(record.last_transport_change_at_ms)
>= MANAGED_SETTLE_DEADLINE_MS;
let action = match snapshot.status {
ManagedConnectionHealthStatus::Healthy => ManagedConnectionBridgeAction::Healthy,
ManagedConnectionHealthStatus::AwaitingReplacement => {
if settle_deadline_elapsed {
let timeout_reason = if matches!(
self.session_admission(connection_id),
crate::session_token::SessionAdmission::Pending
) {
"session-admission-timeout"
} else {
"settle-timeout"
};
ManagedConnectionBridgeAction::Retire {
reason: timeout_reason.to_string(),
}
} else if snapshot
.node_id
.as_deref()
.and_then(|value| value.parse::<iroh::EndpointId>().ok())
.is_some()
{
ManagedConnectionBridgeAction::Rebind {
remote_node_id: snapshot.node_id,
transport_generation: snapshot.transport_generation,
}
} else {
ManagedConnectionBridgeAction::AwaitReplacement {
reason: snapshot.readiness_reason.clone(),
}
}
}
ManagedConnectionHealthStatus::Dead => {
if settle_deadline_elapsed {
ManagedConnectionBridgeAction::Retire {
reason: "settle-timeout".to_string(),
}
} else {
ManagedConnectionBridgeAction::Retire {
reason: "managed transport no longer alive".to_string(),
}
}
}
};
Some(action)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn send_message(
&self,
target_id: &str,
payload: &str,
state: Option<&str>,
reply_payload: Option<&str>,
) -> anyhow::Result<String> {
let node_guard = self.node_id.read().await;
if let Some(node_id) = node_guard.as_deref() {
self.signaling
.send_message(node_id, target_id, payload, state, reply_payload)
.await
} else {
Err(anyhow::anyhow!("Node ID not set"))
}
}
#[cfg(target_arch = "wasm32")]
pub async fn send_message(
&self,
_target_id: &str,
_payload: &str,
_state: Option<&str>,
_reply_payload: Option<&str>,
) -> anyhow::Result<String> {
Err(Self::deprecated_browser_signaling_error("send_message"))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn subscribe_devices(
&self,
user_id: &str,
) -> anyhow::Result<
futures::stream::BoxStream<'static, anyhow::Result<Vec<crate::signaling::DeviceEvent>>>,
> {
self.signaling.subscribe_devices(user_id).await
}
#[cfg(target_arch = "wasm32")]
pub async fn subscribe_devices(
&self,
_user_id: &str,
) -> anyhow::Result<
futures::stream::BoxStream<'static, anyhow::Result<Vec<crate::signaling::DeviceEvent>>>,
> {
Err(Self::deprecated_browser_signaling_error(
"subscribe_devices",
))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn create_session(
&self,
mut session: crate::signaling::SignalingSession,
) -> anyhow::Result<()> {
if session.app_tag.is_none() {
session.app_tag = Some(self.app_tag.clone());
}
self.signaling.create_session(session).await
}
#[cfg(target_arch = "wasm32")]
pub async fn create_session(
&self,
_session: crate::signaling::SignalingSession,
) -> anyhow::Result<()> {
Err(Self::deprecated_browser_signaling_error("create_session"))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn update_session(
&self,
session_id: &str,
update: serde_json::Value,
) -> anyhow::Result<()> {
self.signaling.update_session(session_id, update).await
}
#[cfg(target_arch = "wasm32")]
pub async fn update_session(
&self,
_session_id: &str,
_update: serde_json::Value,
) -> anyhow::Result<()> {
Err(Self::deprecated_browser_signaling_error("update_session"))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn subscribe_sessions(
&self,
local_device_id: &str,
) -> anyhow::Result<
futures::stream::BoxStream<'static, anyhow::Result<Vec<crate::signaling::SessionEvent>>>,
> {
self.signaling.subscribe_sessions(local_device_id).await
}
#[cfg(target_arch = "wasm32")]
pub async fn subscribe_sessions(
&self,
_local_device_id: &str,
) -> anyhow::Result<
futures::stream::BoxStream<'static, anyhow::Result<Vec<crate::signaling::SessionEvent>>>,
> {
Err(Self::deprecated_browser_signaling_error(
"subscribe_sessions",
))
}
#[cfg(not(target_arch = "wasm32"))]
pub fn start_signaling_loop(
self: std::sync::Arc<Self>,
user_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
) -> tokio::task::JoinHandle<()> {
let (tx, mut rx) = tokio::sync::mpsc::channel(8);
{
let mut guard = self.presence_loop_tx.lock().unwrap();
*guard = Some(tx);
}
tokio::spawn(async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(240));
if self.is_app_backgrounded() {
return;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
let (iroh_fingerprint, scope, token_fingerprint) =
super::core_impl::summarize_compound_ticket_for_logs(effective_ticket.as_str());
println!(
"[pluto-rtc][presence][publish] reason=initial scope={} token_fp={} iroh_fp={}",
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string()),
iroh_fingerprint
);
if let Err(e) = self
.update_presence(
&user_id,
&device_name,
&effective_ticket,
metadata.as_deref(),
)
.await
{
eprintln!("Failed to update native signaling presence: {}", e);
}
loop {
tokio::select! {
_ = interval.tick() => {
if self.is_app_backgrounded() {
continue;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
let (iroh_fingerprint, scope, token_fingerprint) =
super::core_impl::summarize_compound_ticket_for_logs(effective_ticket.as_str());
println!(
"[pluto-rtc][presence][publish] reason=heartbeat scope={} token_fp={} iroh_fp={}",
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string()),
iroh_fingerprint
);
if let Err(e) = self
.update_presence(
&user_id,
&device_name,
&effective_ticket,
metadata.as_deref(),
)
.await
{
eprintln!("Failed to update native signaling presence: {}", e);
}
}
cmd = rx.recv() => {
match cmd {
Some(crate::presence::PresenceCommand::UpdateNow) => {
if self.is_app_backgrounded() {
continue;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
let (iroh_fingerprint, scope, token_fingerprint) =
super::core_impl::summarize_compound_ticket_for_logs(effective_ticket.as_str());
println!(
"[pluto-rtc][presence][publish] reason=update-now scope={} token_fp={} iroh_fp={}",
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string()),
iroh_fingerprint
);
if let Err(e) = self
.update_presence(
&user_id,
&device_name,
&effective_ticket,
metadata.as_deref(),
)
.await
{
eprintln!("Failed to force update native signaling presence: {}", e);
}
interval.reset();
}
Some(crate::presence::PresenceCommand::Stop) | None => {
break;
}
}
}
}
}
})
}
#[cfg(target_arch = "wasm32")]
pub fn start_signaling_loop(
self: std::sync::Arc<Self>,
user_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
) {
let (tx, mut rx) = tokio::sync::mpsc::channel(8);
{
let mut guard = self.presence_loop_tx.lock().unwrap();
*guard = Some(tx);
}
wasm_bindgen_futures::spawn_local(async move {
if self.is_app_backgrounded() {
return;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
if let Err(e) = self
.update_presence(
&user_id,
&device_name,
&effective_ticket,
metadata.as_deref(),
)
.await
{
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"Failed to update native signaling presence (initial): {}",
e
)));
}
loop {
let sleep = gloo_timers::future::sleep(std::time::Duration::from_secs(240));
tokio::select! {
_ = sleep => {
if self.is_app_backgrounded() {
continue;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
if let Err(e) = self
.update_presence(&user_id, &device_name, &effective_ticket, metadata.as_deref())
.await
{
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"Failed to update native signaling presence: {}",
e
)));
}
}
cmd = rx.recv() => {
match cmd {
Some(crate::presence::PresenceCommand::UpdateNow) => {
if self.is_app_backgrounded() {
continue;
}
let effective_ticket = self.refresh_presence_ticket_for_publication(&ticket).await;
if let Err(e) = self
.update_presence(&user_id, &device_name, &effective_ticket, metadata.as_deref())
.await
{
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"Failed to force update native signaling presence: {}",
e
)));
}
}
Some(crate::presence::PresenceCommand::Stop) | None => {
break;
}
}
}
}
}
});
}
#[cfg(not(target_arch = "wasm32"))]
pub fn start_heartbeat_loop(
self: std::sync::Arc<Self>,
room_id: String,
member_id: String,
tag: String,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(300)); loop {
interval.tick().await;
if let Err(e) = self.room.heartbeat_tick(&room_id, &member_id, &tag).await {
eprintln!("Failed native room heartbeat tick: {}", e);
}
}
})
}
#[cfg(not(target_arch = "wasm32"))]
pub fn start_auto_connect(self: Arc<Self>, user_id: String, local_device_id: String) {
let key = (user_id.clone(), local_device_id.clone());
let generation = {
let mut guard = match self.auto_connect_loop_key.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
if guard.as_ref() == Some(&key) {
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] duplicate start ignored user_id={} local_device_id={}",
user_id, local_device_id
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] duplicate start ignored user_id={} local_device_id={}",
user_id, local_device_id
)));
return;
}
*guard = Some(key);
self.auto_connect_generation.fetch_add(1, Ordering::SeqCst) + 1
};
let client_clone = self.clone();
#[cfg(not(target_arch = "wasm32"))]
tokio::spawn(async move {
client_clone
.auto_connect_loop(user_id, local_device_id, generation)
.await;
});
#[cfg(target_arch = "wasm32")]
wasm_bindgen_futures::spawn_local(async move {
client_clone
.auto_connect_loop(user_id, local_device_id, generation)
.await;
});
}
#[cfg(target_arch = "wasm32")]
pub fn start_auto_connect(self: Arc<Self>, user_id: String, local_device_id: String) {
let key = (user_id, local_device_id);
let mut guard = match self.auto_connect_loop_key.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
*guard = Some(key);
self.auto_connect_generation
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
#[cfg(target_arch = "wasm32")]
fn deprecated_browser_signaling_error(operation: &str) -> anyhow::Error {
Self::warn_deprecated_browser_signaling(operation);
anyhow::anyhow!(
"The wasm {} browser Firestore path has been retired. Use the TypeScript Firebase signaling adapter instead.",
operation
)
}
#[cfg(target_arch = "wasm32")]
fn warn_deprecated_browser_signaling(operation: &str) {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[OPENRTC][WASM-SIGNALING] {} is deprecated in browsers. Use the TypeScript Firebase signaling path instead.",
operation
)));
}
pub fn force_reconnect_snapshot(self: Arc<Self>) {
if self.is_app_backgrounded() {
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(
"[Client] force_reconnect_snapshot skipped while app is backgrounded",
));
#[cfg(not(target_arch = "wasm32"))]
println!("[Client] force_reconnect_snapshot skipped while app is backgrounded");
return;
}
if let Some(tx) = self.presence_loop_tx.lock().unwrap().clone() {
let _ = tx.try_send(crate::presence::PresenceCommand::UpdateNow);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(
"[Client] Forced presence update triggered",
));
#[cfg(not(target_arch = "wasm32"))]
println!("[Client] Forced presence update triggered");
}
#[cfg(target_arch = "wasm32")]
{
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(
"[Client] force_reconnect_snapshot skipped browser auto-connect loop; TypeScript Firebase signaling owns device snapshots",
));
return;
}
#[cfg(not(target_arch = "wasm32"))]
{
let (user_id, local_device_id) = {
let guard = match self.auto_connect_loop_key.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
match guard.as_ref() {
Some((uid, did)) => (uid.clone(), did.clone()),
None => return, }
};
let generation = self
.auto_connect_generation
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1;
let client_clone = self.clone();
tokio::spawn(async move {
client_clone
.auto_connect_loop(user_id, local_device_id, generation)
.await;
});
}
}
pub fn stop_presence_loop(&self) {
let tx = {
let mut guard = self.presence_loop_tx.lock().unwrap();
guard.take()
};
if let Some(tx) = tx {
let _ = tx.try_send(crate::presence::PresenceCommand::Stop);
}
}
pub fn stop_auto_connect(&self) {
let mut guard = match self.auto_connect_loop_key.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
*guard = None;
self.auto_connect_generation
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
pub fn stop_auth_scoped_activity(&self) {
self.stop_presence_loop();
self.stop_auto_connect();
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn is_auto_connect_generation_current(&self, generation: u64) -> bool {
self.auto_connect_generation.load(Ordering::SeqCst) == generation
}
pub fn set_auto_connect_excluded(&self, device_id: &str, excluded: bool) {
self.set_auto_connect_excluded_peer(device_id, None, excluded);
}
pub(crate) fn set_auto_connect_excluded_peer(
&self,
device_id: &str,
node_id: Option<&str>,
excluded: bool,
) {
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return;
};
let normalized_node_id = node_id.and_then(Self::normalize_auto_connect_exclusion_key);
let mut guard = match self.auto_connect_excluded.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
let mut aliases = match self.auto_connect_excluded_node_aliases.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
if excluded {
guard.insert(device_id.clone());
if let Some(node_id) = normalized_node_id {
aliases.entry(device_id).or_default().insert(node_id);
}
} else {
guard.remove(&device_id);
aliases.remove(&device_id);
}
}
pub fn is_auto_connect_excluded(&self, device_id: &str) -> bool {
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return false;
};
match self.auto_connect_excluded.lock() {
Ok(g) => g.contains(&device_id),
Err(p) => p.into_inner().contains(&device_id),
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn is_auto_connect_peer_excluded(
&self,
device_id: &str,
node_id: Option<&str>,
) -> bool {
if self.is_auto_connect_excluded(device_id) {
return true;
}
let Some(node_id) = node_id.and_then(Self::normalize_auto_connect_exclusion_key) else {
return false;
};
let aliases = match self.auto_connect_excluded_node_aliases.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
aliases.values().any(|values| values.contains(&node_id))
}
pub fn set_app_backgrounded(self: &Arc<Self>, backgrounded: bool) {
let was_backgrounded = self
.app_backgrounded
.swap(backgrounded, std::sync::atomic::Ordering::Relaxed);
if was_backgrounded != backgrounded {
#[cfg(not(target_arch = "wasm32"))]
println!(
"[AppState] set_app_backgrounded transition={} previous={} current={}",
if backgrounded {
"foreground->background"
} else {
"background->foreground"
},
was_backgrounded,
backgrounded
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[AppState] set_app_backgrounded transition={} previous={} current={}",
if backgrounded {
"foreground->background"
} else {
"background->foreground"
},
was_backgrounded,
backgrounded
)));
}
#[cfg(not(target_arch = "wasm32"))]
if was_backgrounded && !backgrounded {
let client = self.clone();
if let Ok(handle) = tokio::runtime::Handle::try_current() {
handle.spawn(async move {
println!("[AppState] Returned to foreground — triggering transport upgrades for relay-connected peers");
for record in client.connection_manager.list_active().await {
let Some(ref endpoint_id) = record.endpoint_id else {
continue;
};
if client.iroh_path_kind(endpoint_id).await
== crate::client::IrohPathKind::Relay
{
if client.is_webrtc_transport_enabled().await {
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
let _ = client
.request_native_webrtc_recovery(
&record.connection_id,
Some(endpoint_id.as_str()),
crate::native_webrtc_policy::NativeWebRTCRecoveryTrigger::Native(
crate::native_webrtc_policy::NativeWebRTCNativeTrigger::ForegroundResume,
),
crate::native_webrtc_policy::NativeWebRTCRecoveryOptions {
force_restart: false,
preferred_negotiation_id: None,
role_override: None,
},
)
.await;
}
}
}
}
});
}
}
}
pub fn is_app_backgrounded(&self) -> bool {
self.app_backgrounded
.load(std::sync::atomic::Ordering::Relaxed)
}
}
#[cfg(test)]
mod tests {
use super::Client;
#[test]
fn normalize_transport_name_accepts_official_path_labels() {
for label in [
"iroh",
"iroh-quic",
"iroh-lan",
"iroh-relay",
"ble",
"webrtc",
"webrtc-lan",
"moq",
] {
assert_eq!(
Client::normalize_transport_name(label).as_deref(),
Some(label),
"label {label} should be accepted"
);
}
}
}