use super::*;
pub(crate) fn should_present_session_token_on_connected_transport(
has_extracted_token: bool,
has_current_remote_admission_proof: bool,
) -> bool {
has_extracted_token && !has_current_remote_admission_proof
}
pub(crate) fn wasm_independent_route_proof_is_current(
active_transport: &str,
health: &crate::connection_manager::ConnectionHealth,
) -> bool {
matches!(health, crate::connection_manager::ConnectionHealth::Healthy)
&& (crate::transport_label::is_webrtc(active_transport)
|| crate::transport_label::normalize(active_transport)
== Some(crate::transport_label::MOQ))
}
fn selected_iroh_latency_label<'a>(
path_kind: crate::client::IrohPathKind,
transport_labels: impl IntoIterator<Item = &'a str>,
) -> Option<&'a str> {
let labels = transport_labels
.into_iter()
.filter_map(|label| crate::transport_label::normalize(label).map(|_| label))
.collect::<Vec<_>>();
let matches_kind = |label: &&str| {
let normalized = crate::transport_label::normalize(label);
match path_kind {
crate::client::IrohPathKind::DirectQuic => {
normalized == Some(crate::transport_label::IROH)
|| normalized == Some(crate::transport_label::IROH_QUIC)
}
crate::client::IrohPathKind::DirectLan => {
normalized == Some(crate::transport_label::IROH_LAN)
}
crate::client::IrohPathKind::Relay => {
normalized == Some(crate::transport_label::IROH_RELAY)
}
crate::client::IrohPathKind::Ble => normalized == Some(crate::transport_label::BLE),
crate::client::IrohPathKind::Unknown => {
normalized.is_some_and(crate::transport_label::is_iroh_base)
}
}
};
labels.into_iter().find(matches_kind)
}
#[cfg(not(target_arch = "wasm32"))]
fn native_rtdb_presence_refresh_ms() -> u64 {
const DEFAULT_MS: u64 = 30_000;
const MIN_MS: u64 = 5_000;
const MAX_MS: u64 = 5 * 60 * 1000;
std::env::var("OPENRTC_RTDB_PRESENCE_REFRESH_MS")
.ok()
.and_then(|value| value.trim().parse::<u64>().ok())
.map(|value| value.clamp(MIN_MS, MAX_MS))
.unwrap_or(DEFAULT_MS)
}
#[cfg(not(target_arch = "wasm32"))]
fn native_presence_retry_targets(durable_registered: bool, live_registered: bool) -> (bool, bool) {
(!durable_registered, !live_registered)
}
#[cfg(not(target_arch = "wasm32"))]
fn native_presence_is_ready(durable_registered: bool, live_registered: bool) -> bool {
durable_registered && live_registered
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Clone, Debug)]
enum NativePresenceTicketPolicy {
Fixed(String),
ManagedUserDevice,
}
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>,
) {
let terminal_disconnect =
crate::lifecycle_reason::reason_is_terminal_disconnect(normalized_reason.as_deref());
self.connection_manager
.set_closing(connection_id, normalized_reason.clone())
.await;
self.connection_manager
.set_closed(connection_id, normalized_reason)
.await;
#[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;
if terminal_disconnect {
self.forget_session_connection(connection_id);
}
let _ = self.connection_manager.remove(connection_id).await;
self.managed_connect_gates
.lock()
.await
.remove(connection_id);
#[cfg(not(target_arch = "wasm32"))]
self.native_peer_transport_capabilities
.write()
.await
.remove(connection_id);
#[cfg(not(target_arch = "wasm32"))]
self.native_control_streams
.lock()
.await
.remove(connection_id);
#[cfg(not(target_arch = "wasm32"))]
self.native_optional_route_generations
.write()
.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_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(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
if let Some(session) = self.native_moq_sessions.write().await.remove(connection_id) {
session.close_gracefully().await;
}
}
#[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)
}
#[cfg(test)]
pub async fn report_transport_status(
&self,
connection_id: &str,
active_transport: &str,
parallel_transport: Option<&str>,
) -> Option<PeerSessionSnapshot> {
self.report_transport_status_for_current_generation(
connection_id,
active_transport,
parallel_transport,
)
.await
}
pub(crate) async fn report_transport_status_for_current_generation(
&self,
connection_id: &str,
active_transport: &str,
parallel_transport: Option<&str>,
) -> Option<PeerSessionSnapshot> {
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await?;
let stable_id = record.transport_stable_id?;
self.report_transport_status_for_generation(
connection_id,
active_transport,
parallel_transport,
stable_id,
record.transport_generation,
record.route_generation,
)
.await
}
pub async fn report_transport_status_for_generation(
&self,
connection_id: &str,
active_transport: &str,
parallel_transport: Option<&str>,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
expected_route_generation: u64,
) -> Option<PeerSessionSnapshot> {
self.report_transport_status_observation(
connection_id,
active_transport,
parallel_transport,
(
expected_transport_stable_id,
expected_transport_generation,
expected_route_generation,
),
)
.await
}
async fn report_transport_status_observation(
&self,
connection_id: &str,
active_transport: &str,
parallel_transport: Option<&str>,
expected_generation: (u64, u64, u64),
) -> 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 (stable_id, transport_generation, route_generation) = expected_generation;
let snapshot = self
.connection_manager
.report_transport_status_if_current(
connection_id,
stable_id,
transport_generation,
route_generation,
active_transport,
parallel_transport,
)
.await
.map(peer_session_snapshot_from_peer)
.map(|snapshot| self.apply_peer_session_admission_guard(snapshot));
#[cfg(not(target_arch = "wasm32"))]
self.refresh_native_optional_route_generation(connection_id)
.await;
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
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn current_native_peer_data_generation(
&self,
connection_id: &str,
expected_transport_stable_id: Option<u64>,
) -> Option<crate::client::NativePeerDataGeneration> {
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await?;
if !matches!(
record.state,
crate::connection_manager::ConnectionState::Connecting
| crate::connection_manager::ConnectionState::Connected
) {
return None;
}
let transport_stable_id = record.transport_stable_id?;
if expected_transport_stable_id.is_some_and(|expected| expected != transport_stable_id) {
return None;
}
Some(crate::client::NativePeerDataGeneration {
transport_stable_id,
transport_generation: record.transport_generation,
route_generation: record.route_generation,
})
}
async fn best_peer_snapshot_for_lookup(
&self,
id: &str,
seed: Option<&crate::connection_manager::PeerSnapshot>,
) -> Option<crate::connection_manager::PeerSnapshot> {
let mut aliases = std::collections::HashSet::new();
if let Some(alias) = normalize_lookup_id(Some(id)) {
aliases.insert(alias);
}
if let Some(seed) = seed {
aliases.extend(peer_snapshot_lookup_aliases(seed));
}
if aliases.is_empty() {
return None;
}
self.connection_manager
.list_peer_snapshots()
.await
.into_iter()
.filter(|snapshot| peer_snapshot_matches_any_alias(snapshot, &aliases))
.fold(None, |current, candidate| {
Some(preferred_peer_snapshot(current, candidate))
})
}
async fn resolved_peer_snapshot_for_lookup(
&self,
id: &str,
) -> Option<crate::connection_manager::PeerSnapshot> {
let snapshot = self.connection_manager.peer_snapshot(id).await;
self.best_peer_snapshot_for_lookup(id, snapshot.as_ref())
.await
.or(snapshot)
}
pub(crate) fn session_admission_block_reason(
&self,
connection_id: &str,
) -> Option<(bool, String)> {
self.session_admission_block_reason_for_transport(connection_id, None)
}
pub(crate) fn session_admission_block_reason_for_transport(
&self,
connection_id: &str,
transport_stable_id: Option<u64>,
) -> Option<(bool, String)> {
#[cfg(target_arch = "wasm32")]
let _ = transport_stable_id;
if self.session_registry_active() {
match self.session_admission(connection_id) {
crate::session_token::SessionAdmission::Accepted { .. } => {
#[cfg(not(target_arch = "wasm32"))]
if !self.native_admission_route_is_ready_for_transport(
connection_id,
transport_stable_id,
) {
if Self::admission_trace_enabled() {
let diagnostic = transport_stable_id
.map(|stable_id| {
self.native_application_stream_pending_diagnostic(
connection_id,
stable_id,
)
})
.unwrap_or_else(|| "transport_stable_id=missing".to_string());
eprintln!(
"[OpenRTC][admission-readiness] connection_id={} reason=native-main-route-pending {}",
connection_id, diagnostic,
);
}
return Some((false, "native-main-route-pending".to_string()));
}
}
crate::session_token::SessionAdmission::Pending => {
return Some((false, "session-admission-pending".to_string()));
}
crate::session_token::SessionAdmission::Rejected { reason } => {
return Some((true, reason));
}
}
}
if self.connection_requires_application_crypto(connection_id)
&& self
.application_crypto_key_for_connection(Some(connection_id))
.is_none()
{
return Some((false, "application-crypto-key-pending".to_string()));
}
#[cfg(not(target_arch = "wasm32"))]
if self.connection_requires_application_crypto_confirmation(connection_id)
&& !self.connection_application_crypto_is_confirmed(connection_id)
{
return Some((false, "application-crypto-confirmation-pending".to_string()));
}
None
}
fn apply_peer_session_admission_guard(
&self,
mut snapshot: PeerSessionSnapshot,
) -> PeerSessionSnapshot {
if let Some(connection_id) = snapshot.active_connection_id.as_deref() {
if let crate::session_token::SessionAdmission::Accepted {
scope: Some(scope), ..
} = self.session_admission(connection_id)
{
let scope = scope.into_inner();
if !scope.trim().is_empty() && !snapshot.scopes.contains(&scope) {
snapshot.scopes.push(scope);
snapshot.scopes.sort();
snapshot.scopes.dedup();
}
}
}
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_for_transport(
&connection_id,
snapshot.active_transport_stable_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_for_transport(
&snapshot.connection_id,
snapshot.active_transport_stable_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_for_transport(
&connection_id,
snapshot.active_transport_stable_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);
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 {
#[cfg(not(target_arch = "wasm32"))]
let native_route_diagnostic = snapshot
.active_connection_id
.as_deref()
.zip(snapshot.active_transport_stable_id)
.map(|(connection_id, transport_stable_id)| {
self.native_application_stream_pending_diagnostic(
connection_id,
transport_stable_id,
)
})
.unwrap_or_else(|| "unavailable".to_string());
#[cfg(target_arch = "wasm32")]
let native_route_diagnostic = "wasm-runtime".to_string();
return Err(anyhow::anyhow!(
"peer {} is not settled for protocol traffic: status={:?} readiness={:?} reason={} transport_stable_id={:?} route=[{}]",
id,
snapshot.status,
snapshot.readiness_state,
snapshot.readiness_reason,
snapshot.active_transport_stable_id,
native_route_diagnostic,
));
}
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 Err(anyhow::anyhow!(
"peer {} has no active transport for protocol traffic",
id
));
}
return Ok((snapshot, endpoint_id));
}
}
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_for_transport(
connection_id,
snapshot.active_transport_stable_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));
}
#[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> {
self.resolved_peer_snapshot_for_lookup(id).await
}
pub async fn peer_session(&self, id: &str) -> Option<PeerSessionSnapshot> {
self.resolved_peer_snapshot_for_lookup(id)
.await
.map(peer_session_snapshot_from_peer)
.map(|snapshot| self.apply_peer_session_admission_guard(snapshot))
}
pub async fn peer_sessions(&self) -> Vec<PeerSessionSnapshot> {
let snapshots: Vec<PeerSessionSnapshot> = self
.connection_manager
.list_peer_snapshots()
.await
.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 {
if snapshot_ref.candidate_connection_ids.is_empty() {
return snapshot;
}
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;
}
}
}
#[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(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;
healthy
}
pub(crate) async fn confirm_managed_connection_readiness_from_transport_proof(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
) -> bool {
let confirmed = self
.connection_manager
.set_health_if_current(
connection_id,
expected_transport_stable_id,
crate::connection_manager::ConnectionHealth::Healthy,
)
.await
.is_some();
#[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;
confirmed
}
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 (send, recv) = self
.open_bi_internal_with_timeout(
endpoint_id,
timeout_ms.map(std::time::Duration::from_millis),
)
.await?;
let key = self
.required_application_crypto_key_for_connection_or_endpoint(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
"outgoing bidirectional application stream",
)
.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,
)
})
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn open_peer_protected_bi(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
let (snapshot, endpoint_id) = self
.resolve_settled_peer_endpoint(id, Some(timeout_ms))
.await?;
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))?;
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
self.set_connection_application_crypto_required(connection_id);
if self
.application_crypto_key_for_connection(Some(connection_id))
.is_none()
|| !self.connection_application_crypto_is_confirmed(connection_id)
{
let public_key = self
.get_or_create_connection_key_agreement(connection_id)
.map_err(|error| {
anyhow::anyhow!(
"application crypto key agreement initialization failed for peer {}: {:?}",
id,
error
)
})?
.public_key_bytes();
self.send_typescript_capability_update(
connection_id,
"protected-application-stream",
Some(public_key),
)
.await?;
}
let started = std::time::Instant::now();
let endpoint_id_text = endpoint_id.to_string();
loop {
let current = self
.peer_session(id)
.await
.ok_or_else(|| anyhow::anyhow!("peer {} disappeared during key agreement", id))?;
let generation_is_current = current.active_connection_id.as_deref()
== Some(connection_id)
&& current.node_id.as_deref() == Some(endpoint_id_text.as_str());
if !generation_is_current {
anyhow::bail!(
"peer {} connection generation changed during application crypto key agreement",
id
);
}
if self
.application_crypto_key_for_connection(Some(connection_id))
.is_some()
&& self.connection_application_crypto_is_confirmed(connection_id)
{
break;
}
if started.elapsed().as_millis() >= timeout_ms as u128 {
anyhow::bail!(
"peer {} application crypto key agreement timed out after {}ms",
id,
timeout_ms
);
}
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
}
let (send, recv) = self
.open_bi_internal_with_timeout(
endpoint_id,
Some(std::time::Duration::from_millis(
timeout_ms
.saturating_sub(
started
.elapsed()
.as_millis()
.try_into()
.unwrap_or(timeout_ms),
)
.max(1),
)),
)
.await?;
let key = self
.required_application_crypto_key_for_connection_or_endpoint(
Some(connection_id),
&endpoint_id,
"outgoing protected bidirectional application stream",
)
.await?
.ok_or_else(|| {
anyhow::anyhow!(
"peer {} completed application crypto negotiation without an installed key",
id
)
})?;
let wrapped = crate::application_crypto_streams::wrap_peer_streams(Some(key), send, recv)
.map_err(|error| {
anyhow::anyhow!("application crypto stream wrap failed: {error}")
})?;
Ok((
Some(connection_id.to_string()),
endpoint_id.to_string(),
wrapped.0,
wrapped.1,
))
}
pub async fn send_peer_application_frame(
&self,
id: &str,
frame: &[u8],
timeout_ms: Option<u64>,
) -> anyhow::Result<()> {
let (_connection_id, _remote_node_id, mut send, _recv) =
self.open_peer_bi(id, timeout_ms).await?;
let envelope = crate::stream_metadata::encode_channel_envelope(
crate::stream_metadata::DEFAULT_PEER_CHANNEL_ID,
None,
)?;
send.write_all(&envelope).await?;
send.write_all(frame).await?;
send.finish_and_wait_for_peer(std::time::Duration::from_secs(2))
.await?;
Ok(())
}
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 (connection_id, remote_node_id, send, _recv) = self
.open_peer_bi_explicit_file_streams(id, timeout_ms)
.await?;
Ok((connection_id, remote_node_id, send))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn open_peer_bi_explicit_file(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
self.open_peer_bi_explicit_file_streams(id, timeout_ms)
.await
}
async fn open_peer_bi_explicit_file_streams(
&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 (send, recv) = self
.open_bi_internal_with_timeout(
endpoint_id,
timeout_ms.map(std::time::Duration::from_millis),
)
.await?;
let connection_id = snapshot.active_connection_id.clone();
let key = self
.required_application_crypto_key_for_connection_or_endpoint(
connection_id.as_deref(),
&endpoint_id,
"outgoing explicit-file application stream",
)
.await?;
let (mut send, recv) = crate::application_crypto_streams::wrap_peer_streams(
key, send, recv,
)
.map_err(|error| anyhow::anyhow!("application crypto stream wrap failed: {error}"))?;
send.write_all(&[crate::explicit_transfer_crypto::EXPLICIT_FILE_PROTOCOL_BYTE])
.await?;
Ok((connection_id, remote_node_id, send, recv))
}
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 (send, recv) = self
.open_bi_internal_with_timeout(
endpoint_id,
timeout_ms.map(std::time::Duration::from_millis),
)
.await?;
let key = self
.required_application_crypto_key_for_connection_or_endpoint(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
"outgoing diagnostic application stream",
)
.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_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 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_for_transport(
connection_id,
snapshot.active_transport_stable_id,
) {
return Err(anyhow::anyhow!(
"peer {} is not admitted for transport-only traffic: {}",
id,
reason
));
}
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
.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?;
match self.open_bi_internal(endpoint_id).await {
Ok((send, recv)) => {
return Ok((
snapshot.active_connection_id,
endpoint_id.to_string(),
send,
recv,
));
}
Err(error) => {
#[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-only open timed out after redial attempts: {}",
id,
error
));
}
}
}
#[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();
let key = self
.required_application_crypto_key_for_connection_or_endpoint(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
"outgoing unidirectional application stream",
)
.await?;
Ok((
connection_id,
remote_node_id,
crate::application_crypto_streams::wrap_peer_send_stream(key, 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> {
if let Some(record) = self.connection_manager.get_by_connection_id(id).await {
if let Some(device_id) = record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
{
return Some(device_id.to_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 uses_transient_identity = admitted_scope.as_ref().is_some_and(|scope| {
super::admission_impl::session_scope_uses_transient_peer_identity(scope.as_str())
});
let resolved_device_id = if uses_transient_identity {
Some(node_id.clone())
} else {
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
}
})
};
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()
&& self.native_admission_route_is_ready_for_transport(
&record.connection_id,
record.transport_stable_id,
),
})
}
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() {
let admission = self.session_admission(connection_id);
let durable_admission = matches!(
&admission,
crate::session_token::SessionAdmission::Accepted {
scope: Some(scope),
..
} if super::admission_impl::session_scope_allows_device_binding(scope.as_str())
);
let legacy_trusted_binding = !self.session_registry_active()
&& matches!(admission, crate::session_token::SessionAdmission::Pending);
if durable_admission || legacy_trusted_binding {
let admitted = self
.mark_trusted_user_device_connection_admitted(connection_id, trimmed)
.await;
if admitted {
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.reconcile_authoritative_device_node(trimmed, node_id)
.await;
}
}
}
} else if matches!(admission, crate::session_token::SessionAdmission::Pending) {
self.connection_manager
.set_device_id(connection_id, trimmed.to_string())
.await;
}
}
self.connection_manager.peer_snapshot(connection_id).await
}
pub async fn bind_node_device_id(&self, node_id: &str, device_id: &str) {
let node_id = node_id.trim();
let device_id = device_id.trim();
if node_id.is_empty() || device_id.is_empty() {
return;
}
self.reconcile_authoritative_device_node(device_id, node_id)
.await;
let records = self.connection_manager.get_by_node_id(node_id).await;
for record in records {
let _ = self
.bind_connection_device_id(&record.connection_id, device_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;
}
if settled {
let _ = self
.connection_manager
.clear_status_reason(
connection_id,
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS),
)
.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 report_managed_connection_settled_for_transport(
&self,
connection_id: &str,
settled: bool,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
expected_route_generation: u64,
) -> Option<crate::connection_manager::PeerSnapshot> {
self.connection_manager
.set_settled_if_current(
connection_id,
expected_transport_stable_id,
expected_transport_generation,
expected_route_generation,
settled,
)
.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;
};
let Some(active_connection_id) = snapshot.connection_ids.first().cloned() else {
return false;
};
let Some(expected_transport_stable_id) = snapshot.active_transport_stable_id else {
return false;
};
#[cfg(target_arch = "wasm32")]
if wasm_independent_route_proof_is_current(
snapshot.active_transport.as_str(),
&snapshot.health,
) {
return true;
}
#[cfg(not(target_arch = "wasm32"))]
if self
.native_active_independent_route_ready(&active_connection_id)
.await
{
return self
.connection_manager
.set_health_if_current(
&active_connection_id,
expected_transport_stable_id,
crate::connection_manager::ConnectionHealth::Healthy,
)
.await
.is_some();
}
let node = self.iroh_node.read().await.as_ref().cloned();
let probe = match node {
Some(node) => {
node.probe_connection(endpoint_id, std::time::Duration::from_secs(4))
.await
}
None => None,
};
let (health, healthy) = match probe {
Some(probe)
if probe.transport_stable_id == expected_transport_stable_id
&& probe.responsive =>
{
(crate::connection_manager::ConnectionHealth::Healthy, true)
}
Some(probe) if probe.transport_stable_id != expected_transport_stable_id => {
(crate::connection_manager::ConnectionHealth::Unknown, false)
}
_ => (crate::connection_manager::ConnectionHealth::Stale, false),
};
let updated = self
.connection_manager
.set_health_if_current(&active_connection_id, expected_transport_stable_id, health)
.await;
healthy && updated.is_some()
}
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<()> {
self.publish_presence_record_with_ttl(user_id, device_name, ticket, true, ttl_ms, metadata)
.await
}
pub async fn update_durable_device_record_with_ttl(
&self,
user_id: &str,
device_name: &str,
ticket: &str,
ttl_ms: u64,
metadata: Option<&str>,
) -> anyhow::Result<()> {
self.publish_presence_record_with_ttl(user_id, device_name, ticket, false, ttl_ms, metadata)
.await
}
pub async fn update_live_presence_record(
&self,
user_id: &str,
device_name: &str,
ticket: &str,
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);
let ticket_node_id = parse_endpoint_ticket(iroh_ticket)?.id.to_string();
#[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();
self.signaling
.update_live_presence(
user_id,
&ticket_node_id,
&effective_ticket,
device_name,
metadata,
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
async fn publish_native_user_device_presence_once(
&self,
reason: &str,
user_id: &str,
device_name: &str,
ticket: &str,
durable_ttl_ms: u64,
metadata: Option<&str>,
publish_durable: bool,
publish_live: bool,
) -> (bool, bool) {
if self.is_app_backgrounded() {
return (false, false);
}
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());
let durable_publish = async {
if !publish_durable {
return false;
}
println!(
"[pluto-rtc][presence][durable-register] reason={} liveness_source=rtdb firestore_online=false scope={} token_fp={} iroh_fp={}",
reason,
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string()),
iroh_fingerprint
);
match self
.update_durable_device_record_with_ttl(
user_id,
device_name,
&effective_ticket,
durable_ttl_ms,
metadata,
)
.await
{
Ok(()) => {
println!(
"[pluto-rtc][presence][durable-register][ok] reason={} liveness_source=rtdb firestore_online=false",
reason
);
true
}
Err(e) => {
eprintln!(
"[pluto-rtc][presence][durable-register][retry] reason={} error={}",
reason, e
);
false
}
}
};
let live_publish = async {
if !publish_live {
return false;
}
match self
.update_live_presence_record(user_id, device_name, &effective_ticket, metadata)
.await
{
Ok(()) => {
println!("[pluto-rtc][presence][rtdb-live][ok] reason={}", reason);
true
}
Err(e) => {
eprintln!(
"[pluto-rtc][presence][rtdb-live][retry] reason={} error={}",
reason, e
);
false
}
}
};
let (durable_registered, live_registered) = tokio::join!(durable_publish, live_publish);
(durable_registered, live_registered)
}
#[cfg(not(target_arch = "wasm32"))]
async fn resolve_native_presence_ticket(
&self,
policy: &NativePresenceTicketPolicy,
) -> anyhow::Result<String> {
match policy {
NativePresenceTicketPolicy::Fixed(ticket) => {
Ok(self.refresh_presence_ticket_for_publication(ticket).await)
}
NativePresenceTicketPolicy::ManagedUserDevice => {
let ticket = self.endpoint_ticket_with_token("user-device", 0).await?;
let (iroh_ticket, suffix) = crate::session_token::split_compound_ticket(&ticket);
let payload = suffix
.and_then(|value| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, value)
})
.ok_or_else(|| {
anyhow::anyhow!(
"managed user-device presence requires a compound admission ticket"
)
})?;
if payload.scope.as_str() != "user-device" {
return Err(anyhow::anyhow!(
"managed user-device presence minted unexpected scope {}",
payload.scope
));
}
Ok(ticket)
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn publish_native_presence_with_policy_once(
&self,
reason: &str,
user_id: &str,
device_name: &str,
policy: &NativePresenceTicketPolicy,
durable_ttl_ms: u64,
metadata: Option<&str>,
publish_durable: bool,
publish_live: bool,
) -> (bool, bool) {
let ticket = match self.resolve_native_presence_ticket(policy).await {
Ok(ticket) => ticket,
Err(error) => {
eprintln!(
"[pluto-rtc][presence][ticket-unavailable] reason={} policy={} error={}",
reason,
match policy {
NativePresenceTicketPolicy::Fixed(_) => "fixed",
NativePresenceTicketPolicy::ManagedUserDevice => "managed-user-device",
},
error
);
return (false, false);
}
};
self.publish_native_user_device_presence_once(
reason,
user_id,
device_name,
&ticket,
durable_ttl_ms,
metadata,
publish_durable,
publish_live,
)
.await
}
async fn publish_presence_record_with_ttl(
&self,
user_id: &str,
device_name: &str,
ticket: &str,
is_online: bool,
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);
let ticket_node_id = parse_endpoint_ticket(iroh_ticket)?.id.to_string();
#[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();
self.signaling
.update_presence(
user_id,
&ticket_node_id,
&effective_ticket,
is_online,
device_name,
ttl_ms,
metadata,
)
.await
}
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<()> {
self.stop_presence_loop();
if let Some(node_id) = self.current_node_id().await {
self.signaling
.set_live_presence_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<()> {
if let Some(node_id) = self.current_node_id().await {
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, remote_device_id: &str) {
self.set_auto_connect_excluded_peer(remote_device_id, None, true);
#[cfg(target_arch = "wasm32")]
return;
#[cfg(not(target_arch = "wasm32"))]
let Some((user_id, _)) = self.active_session_identity() else {
return;
};
#[cfg(not(target_arch = "wasm32"))]
let excluded = self.current_excluded_peers_snapshot();
#[cfg(not(target_arch = "wasm32"))]
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, remote_device_id: &str) {
self.set_auto_connect_excluded_peer(remote_device_id, None, false);
#[cfg(target_arch = "wasm32")]
return;
#[cfg(not(target_arch = "wasm32"))]
let Some((user_id, _)) = self.active_session_identity() else {
return;
};
#[cfg(not(target_arch = "wasm32"))]
let excluded = self.current_excluded_peers_snapshot();
#[cfg(not(target_arch = "wasm32"))]
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_id = self.current_node_id().await;
self.signaling
.search_devices(user_id, node_id.as_deref())
.await
}
pub async fn devices_with_status(
&self,
user_id: &str,
) -> anyhow::Result<Vec<DeviceStatusSnapshot>> {
let node_id = self.current_node_id().await;
let devices = self
.signaling
.list_devices(user_id, node_id.as_deref())
.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);
let mut snapshots = Vec::with_capacity(merged.len());
for snapshot in merged {
let mut snapshot = self.apply_device_status_admission_guard(snapshot);
self.enrich_device_status_latency(&mut snapshot).await;
snapshots.push(snapshot);
}
Ok(snapshots)
}
async fn enrich_device_status_latency(&self, snapshot: &mut DeviceStatusSnapshot) {
snapshot.latency_ms = None;
snapshot.latency_by_transport = crate::client::TransportLatencySnapshot::default();
if !snapshot.settled_ready {
return;
}
let expected_generation = match snapshot.connection_id.as_deref() {
Some(connection_id) => {
self.connection_manager
.get_by_connection_id(connection_id)
.await
}
None => None,
};
if expected_generation.as_ref().is_some_and(|record| {
record.transport_stable_id != snapshot.active_transport_stable_id
|| record.transport_generation != snapshot.transport_generation
|| record.route_generation != snapshot.route_generation
}) {
return;
}
let transport_labels = [
Some(snapshot.active_transport.as_str()),
snapshot.parallel_transport.as_deref(),
];
let lookup_id = snapshot
.connection_id
.as_deref()
.or(snapshot.peer_id.as_deref())
.or(snapshot.device.node_id.as_deref())
.unwrap_or(snapshot.device.device_id.as_str());
let mut latency_by_transport = crate::client::TransportLatencySnapshot::default();
if transport_labels.iter().flatten().any(|transport| {
crate::transport_label::normalize(transport)
.is_some_and(crate::transport_label::is_iroh_base)
|| crate::transport_label::normalize(transport) == Some(crate::transport_label::BLE)
}) {
if let Some(latency_ms) = self.iroh_transport_rtt_ms(lookup_id).await {
let path_kind = self.iroh_path_kind(lookup_id).await;
if let Some(transport) = selected_iroh_latency_label(
path_kind,
transport_labels.iter().flatten().copied(),
) {
latency_by_transport.set(transport, latency_ms);
}
}
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
if transport_labels.iter().flatten().any(|transport| {
matches!(
transport.trim().to_ascii_lowercase().as_str(),
"webrtc" | "webrtc-lan" | "webrtc-turn"
)
}) {
for connection_id in self.resolve_transport_connection_ids(lookup_id).await {
let session = self
.native_webrtc_sessions
.read()
.await
.get(&connection_id)
.cloned();
let Some(session) = session else { continue };
if !self
.current_native_webrtc_route_identity_matches(&connection_id, &session)
.await
{
continue;
}
if let Some(latency_ms) = session.latency_ms().await {
for transport in transport_labels.iter().flatten().copied() {
if matches!(
transport.trim().to_ascii_lowercase().as_str(),
"webrtc" | "webrtc-lan" | "webrtc-turn"
) {
latency_by_transport.set(transport, latency_ms);
}
}
break;
}
}
}
if let Some(expected) = expected_generation {
let current = self
.connection_manager
.get_by_connection_id(&expected.connection_id)
.await;
if current.as_ref().is_none_or(|record| {
record.transport_stable_id != expected.transport_stable_id
|| record.transport_generation != expected.transport_generation
|| record.route_generation != expected.route_generation
}) {
return;
}
}
snapshot.latency_by_transport = latency_by_transport;
snapshot.latency_ms = snapshot
.latency_by_transport
.get(&snapshot.active_transport)
.or_else(|| {
snapshot
.parallel_transport
.as_deref()
.and_then(|transport| snapshot.latency_by_transport.get(transport))
});
}
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;
#[cfg(not(target_arch = "wasm32"))]
let ble_available = self
.native_custom_transport_kinds
.read()
.await
.values()
.any(|kind| matches!(kind, IrohPathKind::Ble));
#[cfg(target_arch = "wasm32")]
let ble_available = false;
let transport = RuntimeTransportStatus::from_config_with_ble_available(
&transport_config,
ble_available,
current_iroh_relay_provider(),
);
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);
}
self.wake_native_ble_recovery("network-change").await;
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> {
self.connect_device_with_intent(device_id, endpoint_ticket, true)
.await
}
#[cfg(target_arch = "wasm32")]
pub(crate) async fn connect_desired_device(
&self,
device_id: &str,
endpoint_ticket: &str,
) -> anyhow::Result<ManagedConnectResult> {
self.connect_device_with_intent(Some(device_id), endpoint_ticket, false)
.await
}
async fn connect_device_with_intent(
&self,
device_id: Option<&str>,
endpoint_ticket: &str,
clear_auto_connect_exclusion: bool,
) -> 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)?;
self.remember_endpoint_addr(&endpoint_addr).await;
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);
#[cfg(target_arch = "wasm32")]
self.set_connection_application_crypto_required(&connection_id);
let trimmed_device_id = device_id
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned);
let ensure_desired_node_is_current = || -> anyhow::Result<()> {
if clear_auto_connect_exclusion {
return Ok(());
}
let Some(device_id) = trimmed_device_id.as_deref() else {
return Ok(());
};
let Some(authoritative_node_id) = self.authoritative_node_for_device(device_id) else {
return Ok(());
};
if authoritative_node_id == remote_node_id {
return Ok(());
}
Err(anyhow::anyhow!(
"browser desired-peer generation is stale: device {device_id} now belongs to node {authoritative_node_id}, not {remote_node_id}"
))
};
if clear_auto_connect_exclusion {
if let Some(device_id) = trimmed_device_id.as_deref() {
self.unexclude_peer_and_publish(device_id).await;
}
}
let connect_gate = self.managed_connect_gate(&connection_id).await;
let _connect_guard = connect_gate.lock().await;
ensure_desired_node_is_current()?;
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
{
let live_transport_stable_id = self
.get_connection(endpoint_addr.id)
.await
.map(|connection| connection.stable_id() as u64);
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.clone()),
live_transport_stable_id,
Some("connect-device-existing-live-transport".to_string()),
)
.await;
if let Some(device_id) = trimmed_device_id.clone() {
self.mark_trusted_user_device_connection_admitted(&connection_id, &device_id)
.await;
self.reconcile_authoritative_device_node(&device_id, &remote_node_id)
.await;
}
let has_remote_admission_proof = if let Some(token) = extracted_token.as_deref() {
self.has_current_remote_session_admission_proof(
&connection_id,
endpoint_addr.id,
token,
)
.await
} else {
false
};
let needs_token_representation =
should_present_session_token_on_connected_transport(
extracted_token.is_some(),
has_remote_admission_proof,
);
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_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await
{
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
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
);
}
}
ensure_desired_node_is_current()?;
}
}
ensure_desired_node_is_current()?;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
ensure_desired_node_is_current()?;
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.reconcile_authoritative_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?;
ensure_desired_node_is_current()?;
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;
ensure_desired_node_is_current()?;
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;
ensure_desired_node_is_current()?;
log_phase("trusted-admission-complete");
log_phase("stale-record-prune-start");
self.reconcile_authoritative_device_node(&device_id, &remote_node_id)
.await;
log_phase("stale-record-prune-complete");
}
if let Some(ref token) = extracted_token {
match self
.present_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await
{
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
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;
ensure_desired_node_is_current()?;
}
ensure_desired_node_is_current()?;
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?;
let 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_for_transport(
connection_id,
record.transport_stable_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 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_state_change_at_ms
.max(record.last_transport_change_at_ms)
.max(record.last_route_change_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,
route_generation: record.route_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 (handle, _readiness) = self.start_native_presence_loop(
user_id,
device_name,
NativePresenceTicketPolicy::Fixed(ticket),
metadata,
);
handle
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn start_managed_user_device_presence_loop(
self: std::sync::Arc<Self>,
user_id: String,
device_name: String,
metadata: Option<String>,
) -> anyhow::Result<tokio::task::JoinHandle<()>> {
const INITIAL_PRESENCE_READY_TIMEOUT_SECS: u64 = 15;
let owner = self.clone();
let (handle, readiness) = self.start_native_presence_loop(
user_id,
device_name,
NativePresenceTicketPolicy::ManagedUserDevice,
metadata,
);
match tokio::time::timeout(
std::time::Duration::from_secs(INITIAL_PRESENCE_READY_TIMEOUT_SECS),
readiness,
)
.await
{
Ok(Ok(())) => Ok(handle),
Ok(Err(_)) => {
owner.stop_presence_loop();
handle.abort();
Err(anyhow::anyhow!(
"native presence actor exited before initial publication completed"
))
}
Err(_) => {
owner.stop_presence_loop();
handle.abort();
Err(anyhow::anyhow!(
"native presence did not publish durable and RTDB records within {} seconds",
INITIAL_PRESENCE_READY_TIMEOUT_SECS
))
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn start_native_presence_loop(
self: std::sync::Arc<Self>,
user_id: String,
device_name: String,
ticket_policy: NativePresenceTicketPolicy,
metadata: Option<String>,
) -> (
tokio::task::JoinHandle<()>,
tokio::sync::oneshot::Receiver<()>,
) {
let (tx, mut rx) = tokio::sync::mpsc::channel(8);
let actor_tx = tx.clone();
let (readiness_tx, readiness_rx) = tokio::sync::oneshot::channel();
let replaced_sender = {
let mut guard = self.presence_loop_tx.lock().unwrap();
guard.replace(tx)
};
let replaced_existing_actor = replaced_sender.is_some();
if let Some(replaced_sender) = replaced_sender {
let _ = replaced_sender.try_send(crate::presence::PresenceCommand::Stop);
}
let handle = tokio::spawn(async move {
const DURABLE_DEVICE_RECORD_TTL_MS: u64 = 7 * 24 * 60 * 60 * 1000;
const NATIVE_PRESENCE_INITIAL_RETRY_MS: u64 = 2_000;
let mut durable_registered = false;
let mut live_registered = false;
let mut readiness_tx = Some(readiness_tx);
let policy_label = match &ticket_policy {
NativePresenceTicketPolicy::Fixed(_) => "fixed",
NativePresenceTicketPolicy::ManagedUserDevice => "managed-user-device",
};
eprintln!(
"[pluto-rtc][presence][actor-start] user_id={} device_name={} policy={} replaced_existing_actor={}",
user_id, device_name, policy_label, replaced_existing_actor
);
let mut retry_interval = tokio::time::interval(std::time::Duration::from_millis(
NATIVE_PRESENCE_INITIAL_RETRY_MS,
));
retry_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut live_refresh_interval = tokio::time::interval(
std::time::Duration::from_millis(native_rtdb_presence_refresh_ms()),
);
live_refresh_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = retry_interval.tick(), if !durable_registered || !live_registered => {
let (publish_durable, publish_live) =
native_presence_retry_targets(durable_registered, live_registered);
let (durable_ok, live_ok) = self
.publish_native_presence_with_policy_once(
if durable_registered || live_registered {
"initial-retry"
} else {
"initial"
},
&user_id,
&device_name,
&ticket_policy,
DURABLE_DEVICE_RECORD_TTL_MS,
metadata.as_deref(),
publish_durable,
publish_live,
)
.await;
durable_registered = durable_registered || durable_ok;
live_registered = live_registered || live_ok;
if native_presence_is_ready(durable_registered, live_registered) {
if let Some(readiness_tx) = readiness_tx.take() {
let _ = readiness_tx.send(());
eprintln!(
"[pluto-rtc][presence][actor-ready] user_id={} device_name={} policy={}",
user_id, device_name, policy_label
);
}
}
}
_ = live_refresh_interval.tick(), if live_registered => {
let refresh = match self.resolve_native_presence_ticket(&ticket_policy).await {
Ok(ticket) => self
.update_live_presence_record(
&user_id,
&device_name,
&ticket,
metadata.as_deref(),
)
.await,
Err(error) => Err(error),
};
match refresh {
Ok(()) => {
println!("[pluto-rtc][presence][rtdb-live][ok] reason=lease-refresh");
}
Err(error) => {
live_registered = false;
eprintln!(
"[pluto-rtc][presence][rtdb-live][retry] reason=lease-refresh error={}",
error
);
}
}
}
cmd = rx.recv() => {
match cmd {
Some(crate::presence::PresenceCommand::RefreshLiveNow) => {
let (_, live_ok) = self
.publish_native_presence_with_policy_once(
"live-refresh-now",
&user_id,
&device_name,
&ticket_policy,
DURABLE_DEVICE_RECORD_TTL_MS,
metadata.as_deref(),
false,
true,
)
.await;
live_registered = live_ok;
}
Some(crate::presence::PresenceCommand::RepublishDurableNow) => {
let (durable_ok, live_ok) = self
.publish_native_presence_with_policy_once(
"durable-republish-now",
&user_id,
&device_name,
&ticket_policy,
DURABLE_DEVICE_RECORD_TTL_MS,
metadata.as_deref(),
true,
true,
)
.await;
durable_registered = durable_ok;
live_registered = live_ok;
}
Some(crate::presence::PresenceCommand::Stop) => {
eprintln!(
"[pluto-rtc][presence][actor-stop] user_id={} device_name={} reason=command",
user_id, device_name
);
break;
}
None => {
eprintln!(
"[pluto-rtc][presence][actor-stop] user_id={} device_name={} reason=sender-replaced-or-dropped",
user_id, device_name
);
break;
}
}
}
}
}
let mut guard = self.presence_loop_tx.lock().unwrap();
if guard
.as_ref()
.is_some_and(|active_tx| active_tx.same_channel(&actor_tx))
{
guard.take();
}
});
(handle, readiness_rx)
}
#[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 {
loop {
tokio::select! {
cmd = rx.recv() => {
match cmd {
Some(crate::presence::PresenceCommand::RefreshLiveNow)
| Some(crate::presence::PresenceCommand::RepublishDurableNow) => {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(
"[OPENRTC][WASM-SIGNALING] start_presence_loop is disabled. Browser liveness uses the TypeScript RTDB presence adapter; use update_presence only for explicit durable ticket registration.",
));
}
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;
}
self.known_device_ids_by_node
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
self.native_route_repair_credentials
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
*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(),
};
if guard.as_ref() != Some(&key) {
self.known_device_ids_by_node
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
}
*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::RepublishDurableNow);
#[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 request_presence_update(&self) -> bool {
self.presence_loop_tx
.lock()
.unwrap()
.as_ref()
.is_some_and(|tx| {
tx.try_send(crate::presence::PresenceCommand::RefreshLiveNow)
.is_ok()
})
}
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);
#[cfg(not(target_arch = "wasm32"))]
self.native_route_repair_credentials
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
}
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());
match self.auto_connect_peer_requested_excluded.lock() {
Ok(mut peer_requested) => {
peer_requested.remove(&device_id);
}
Err(poisoned) => {
poisoned.into_inner().remove(&device_id);
}
}
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);
match self.auto_connect_peer_requested_excluded.lock() {
Ok(mut peer_requested) => {
peer_requested.remove(&device_id);
}
Err(poisoned) => {
poisoned.into_inner().remove(&device_id);
}
}
}
}
pub(crate) fn set_peer_requested_auto_connect_excluded(
&self,
device_id: &str,
node_id: Option<&str>,
) {
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 excluded = match self.auto_connect_excluded.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
let mut aliases = match self.auto_connect_excluded_node_aliases.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
let mut peer_requested = match self.auto_connect_peer_requested_excluded.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
let locally_excluded =
excluded.contains(&device_id) && !peer_requested.contains(&device_id);
if let Some(node_id) = normalized_node_id {
aliases
.entry(device_id.clone())
.or_default()
.insert(node_id);
}
if locally_excluded {
return;
}
excluded.insert(device_id.clone());
peer_requested.insert(device_id);
}
pub(crate) fn resume_peer_requested_auto_connect_if_authenticated(
&self,
device_id: &str,
) -> bool {
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return false;
};
let peer_requested = match self.auto_connect_peer_requested_excluded.lock() {
Ok(mut peer_requested) => peer_requested.remove(&device_id),
Err(poisoned) => poisoned.into_inner().remove(&device_id),
};
if peer_requested {
self.set_auto_connect_excluded_peer(&device_id, None, false);
}
peer_requested
}
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::{
native_presence_is_ready, native_presence_retry_targets, selected_iroh_latency_label,
Client, NativePresenceTicketPolicy,
};
#[test]
fn native_presence_retry_does_not_repeat_successful_firestore_registration() {
assert_eq!(native_presence_retry_targets(false, false), (true, true));
assert_eq!(native_presence_retry_targets(true, false), (false, true));
assert_eq!(native_presence_retry_targets(false, true), (true, false));
assert_eq!(native_presence_retry_targets(true, true), (false, false));
}
#[test]
fn native_presence_reports_ready_only_after_both_publications() {
assert!(!native_presence_is_ready(false, false));
assert!(!native_presence_is_ready(true, false));
assert!(!native_presence_is_ready(false, true));
assert!(native_presence_is_ready(true, true));
}
#[tokio::test]
async fn managed_native_presence_policy_never_falls_back_to_a_plain_ticket() {
let base_dir = std::env::temp_dir().join(format!(
"openrtc-managed-presence-policy-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("unix epoch")
.as_nanos()
));
let client = Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
"test-app".to_string(),
Box::new(|| None),
);
client
.init_native_device_identity(base_dir.clone(), Some("Managed Presence Device"))
.await
.expect("native identity");
client.init_iroh(None, Vec::new()).await.expect("iroh");
let ticket = client
.resolve_native_presence_ticket(&NativePresenceTicketPolicy::ManagedUserDevice)
.await
.expect("managed ticket");
let (iroh_ticket, suffix) = crate::session_token::split_compound_ticket(&ticket);
let payload = suffix
.and_then(|value| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, value)
})
.expect("compound admission payload");
assert_eq!(payload.scope.as_str(), "user-device");
assert!(!payload.token.trim().is_empty());
let _ = tokio::fs::remove_dir_all(base_dir).await;
}
#[tokio::test]
async fn fixed_native_presence_policy_preserves_explicit_ticket_semantics() {
let client = Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
"test-app".to_string(),
Box::new(|| None),
);
let ticket = client
.resolve_native_presence_ticket(&NativePresenceTicketPolicy::Fixed(
"explicit-ticket".to_string(),
))
.await
.expect("fixed ticket");
assert_eq!(ticket, "explicit-ticket");
}
#[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"
);
}
}
#[test]
fn iroh_latency_is_attributed_only_to_the_selected_physical_path() {
let labels = ["iroh-lan", "iroh-relay"];
assert_eq!(
selected_iroh_latency_label(crate::client::IrohPathKind::DirectLan, labels),
Some("iroh-lan"),
);
assert_eq!(
selected_iroh_latency_label(crate::client::IrohPathKind::Relay, labels),
Some("iroh-relay"),
);
assert_eq!(
selected_iroh_latency_label(crate::client::IrohPathKind::Ble, ["ble", "iroh-relay"],),
Some("ble"),
);
}
}