use super::*;
pub(crate) fn should_preserve_proven_iroh_carrier_route(
current_transport: &str,
observed_transport: &str,
) -> bool {
crate::transport_label::is_generation_bound_iroh_carrier(current_transport)
&& crate::transport_label::is_iroh_base(observed_transport)
}
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
}
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::WebRtc => {
normalized == Some(crate::transport_label::WEBRTC)
}
crate::client::IrohPathKind::Moq => normalized == Some(crate::transport_label::MOQ),
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_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"))]
fn native_presence_retry_delay(consecutive_failures: u32) -> std::time::Duration {
const BASE_MS: u64 = 2_000;
const MAX_MS: u64 = 60_000;
let exponent = consecutive_failures.min(5);
std::time::Duration::from_millis(BASE_MS.saturating_mul(1_u64 << exponent).min(MAX_MS))
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Clone, Debug)]
enum NativePresenceTicketPolicy {
Fixed(String),
ManagedUserDevice,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Default)]
pub(super) struct RetiredNativeCarrierAttempts {
#[cfg(feature = "transport-webrtc")]
webrtc: Option<Arc<crate::client::NativeWebRtcCarrierAttempt>>,
#[cfg(feature = "transport-moq")]
moq: Option<Arc<crate::client::NativeMoqCarrierAttempt>>,
}
#[cfg(not(target_arch = "wasm32"))]
fn native_carrier_generation_matches_retirement(
generation: crate::client::NativePeerDataGeneration,
retiring_generation: Option<(Option<u64>, u64, u64)>,
) -> bool {
retiring_generation.is_some_and(
|(transport_stable_id, transport_generation, route_generation)| {
transport_stable_id.is_none_or(|stable_id| stable_id == generation.transport_stable_id)
&& transport_generation == generation.transport_generation
&& route_generation == generation.route_generation
},
)
}
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()
}
#[cfg(not(target_arch = "wasm32"))]
pub(super) async fn retire_native_carrier_upgrade_state_if_current(
&self,
connection_id: &str,
retiring_generation: Option<(Option<u64>, u64, u64)>,
retiring_peer_policy_epoch: Option<u64>,
) -> RetiredNativeCarrierAttempts {
#[cfg(not(feature = "iroh-carrier-core"))]
let _ = retiring_peer_policy_epoch;
let mut gates = self.native_transport_upgrade_gates.lock().await;
#[cfg(feature = "transport-webrtc")]
let webrtc = {
let mut attempts = self.native_webrtc_carrier_attempts.write().await;
if attempts.get(connection_id).is_some_and(|attempt| {
native_carrier_generation_matches_retirement(
attempt.generation,
retiring_generation,
)
}) {
attempts.remove(connection_id)
} else {
None
}
};
#[cfg(feature = "transport-moq")]
let moq = {
let mut attempts = self.native_moq_carrier_attempts.write().await;
if attempts.get(connection_id).is_some_and(|attempt| {
native_carrier_generation_matches_retirement(
attempt.generation,
retiring_generation,
)
}) {
attempts.remove(connection_id)
} else {
None
}
};
let mut attempts = self.native_ble_upgrade_attempts.lock().await;
if attempts.get(connection_id).is_some_and(|attempt| {
native_carrier_generation_matches_retirement(attempt.generation, retiring_generation)
}) {
attempts.remove(connection_id);
}
drop(attempts);
gates.retain(|(candidate_connection_id, _), gate| {
candidate_connection_id != connection_id
|| !native_carrier_generation_matches_retirement(
gate.generation,
retiring_generation,
)
});
#[cfg(feature = "iroh-carrier-core")]
if !gates
.keys()
.any(|(candidate_connection_id, _)| candidate_connection_id == connection_id)
{
if retiring_peer_policy_epoch.is_some_and(|expected_epoch| {
self.current_iroh_carrier_peer_policy_epoch(connection_id) == expected_epoch
}) {
self.forget_iroh_carrier_peer_policy_epoch(connection_id);
}
}
RetiredNativeCarrierAttempts {
#[cfg(feature = "transport-webrtc")]
webrtc,
#[cfg(feature = "transport-moq")]
moq,
}
}
pub(crate) async fn retire_managed_connection_now(
&self,
connection_id: &str,
normalized_reason: Option<String>,
) {
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "transport-webrtc", feature = "transport-moq")
))]
let retirement_reason = normalized_reason.clone();
#[cfg(not(target_arch = "wasm32"))]
let retiring_record = self
.connection_manager
.get_by_connection_id(connection_id)
.await;
#[cfg(not(target_arch = "wasm32"))]
let retiring_generation = retiring_record.as_ref().map(|record| {
(
record.transport_stable_id,
record.transport_generation,
record.route_generation,
)
});
#[cfg(all(not(target_arch = "wasm32"), feature = "iroh-carrier-core"))]
let retiring_peer_policy_epoch =
Some(self.current_iroh_carrier_peer_policy_epoch(connection_id));
#[cfg(all(not(target_arch = "wasm32"), not(feature = "iroh-carrier-core")))]
let retiring_peer_policy_epoch = None;
let terminal_disconnect =
crate::lifecycle_reason::is_terminal(normalized_reason.as_deref());
#[cfg(not(target_arch = "wasm32"))]
let already_terminal = retiring_record.as_ref().is_some_and(|record| {
matches!(
record.state,
crate::connection_manager::ConnectionState::Closed
| crate::connection_manager::ConnectionState::Failed
)
});
#[cfg(target_arch = "wasm32")]
let already_terminal = self
.connection_manager
.get_by_connection_id(connection_id)
.await
.is_some_and(|record| {
matches!(
record.state,
crate::connection_manager::ConnectionState::Closed
| crate::connection_manager::ConnectionState::Failed
)
});
if !already_terminal {
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(all(target_arch = "wasm32", feature = "iroh-carrier-core"))]
self.forget_iroh_carrier_peer_policy_epoch(connection_id);
#[cfg(not(target_arch = "wasm32"))]
self.native_control_streams
.lock()
.await
.remove(connection_id);
#[cfg(not(target_arch = "wasm32"))]
let retired_native_carrier_attempts = self
.retire_native_carrier_upgrade_state_if_current(
connection_id,
retiring_generation,
retiring_peer_policy_epoch,
)
.await;
#[cfg(all(
not(target_arch = "wasm32"),
not(any(feature = "transport-webrtc", feature = "transport-moq"))
))]
let _ = retired_native_carrier_attempts;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
if let Some(attempt) = retired_native_carrier_attempts.webrtc {
if retirement_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED)
{
if let Some(pump) = attempt.pump.lock().await.as_ref() {
let _ = pump
.send_terminal(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED)
.await;
}
}
self.record_native_iroh_carrier_debug_event(
connection_id,
&attempt.upgrade_id,
"close-requested",
retirement_reason
.as_deref()
.or(Some("logical-connection-retired")),
);
attempt.channel.close();
let _ = attempt.pump.lock().await.take();
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-moq"))]
if let Some(attempt) = retired_native_carrier_attempts.moq {
if retirement_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED)
{
if let Some(pump) = attempt.pump.lock().await.as_ref() {
let _ = pump
.send_terminal(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED)
.await;
}
}
attempt.session.close();
let _ = attempt.pump.lock().await.take();
}
}
#[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
}
#[cfg(any(test, target_arch = "wasm32", feature = "iroh-carrier-core"))]
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(crate) 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;
let (normalized_active_transport, normalized_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 should_preserve_proven_iroh_carrier_route(
record.active_transport.as_str(),
active_transport.as_str(),
) {
active_transport = record.active_transport.clone();
parallel_transport = record.parallel_transport.clone();
}
}
(active_transport, parallel_transport)
};
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));
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,
})
}
#[cfg(all(
target_arch = "wasm32",
any(feature = "transport-webrtc", feature = "transport-moq")
))]
pub(crate) async fn current_wasm_peer_data_generation(
&self,
connection_id: &str,
expected_transport_stable_id: Option<u64>,
) -> Option<crate::client::WasmPeerDataGeneration> {
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::WasmPeerDataGeneration {
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)
}
#[cfg(test)]
pub(crate) fn session_admission_block_reason(
&self,
connection_id: &str,
) -> Option<(bool, String)> {
let block = self.session_admission_block_reason_for_transport(connection_id, None);
if matches!(
block.as_ref(),
Some((false, reason)) if reason == "application-crypto-confirmation-pending"
) && self.connection_application_crypto_has_bound_confirmation(connection_id)
{
None
} else {
block
}
}
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()));
}
if self.connection_requires_application_crypto_confirmation(connection_id)
&& !self.connection_application_crypto_is_confirmed(connection_id, transport_stable_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: StateSnapshot) -> StateSnapshot {
if matches!(
snapshot.state.as_str(),
"closed" | "failed" | "disconnected"
) {
return snapshot;
}
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 {
if matches!(
snapshot.connection_status,
ConnectionStatus::Disconnected | ConnectionStatus::Failed | ConnectionStatus::Closed
) {
return snapshot;
}
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 = ConnectionStatus::Failed;
snapshot.readiness_state = ReadinessState::Failed;
snapshot.peer_health = crate::connection_manager::ConnectionHealth::Stale;
} else {
snapshot.connection_status = ConnectionStatus::Connecting;
snapshot.readiness_state = ReadinessState::Settling;
}
snapshot
}
pub(crate) 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_peer(id, Some(timeout_ms))
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
if snapshot.settled_ready && !Self::peer_session_is_stably_settled(&snapshot) {
anyhow::bail!(
"peer-stream-route-pending: peer {id} generation has not remained settled long enough for protocol traffic"
);
}
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));
}
}
pub(crate) async fn wait_for_peer_transport_replacement(
&self,
id: &str,
prior_connection_id: &str,
prior_endpoint_id: iroh::EndpointId,
prior_transport_stable_id: u64,
timeout_ms: u64,
) -> Option<(PeerSessionSnapshot, iroh::EndpointId)> {
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
if let Some(snapshot) = self.peer_session(id).await {
let endpoint_id = snapshot
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.and_then(|value| value.parse::<iroh::EndpointId>().ok());
let generation_changed = snapshot.active_connection_id.as_deref()
!= Some(prior_connection_id)
|| snapshot.active_transport_stable_id != Some(prior_transport_stable_id)
|| endpoint_id.is_some_and(|value| value != prior_endpoint_id);
if generation_changed && Self::peer_session_is_stably_settled(&snapshot) {
if let Some(endpoint_id) = endpoint_id {
if self.is_connection_transport_alive(endpoint_id).await {
return Some((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 None;
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(50)).await;
}
}
async fn resolve_diagnostic_peer_endpoint(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(PeerSessionSnapshot, iroh::EndpointId)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self
.peer_session(id)
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
let connection_id = snapshot
.active_connection_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("peer {} is missing an active connection", id))?;
let readiness_block = self.session_admission_block_reason_for_transport(
connection_id,
snapshot.active_transport_stable_id,
);
if let Some((true, reason)) = readiness_block.as_ref() {
return Err(anyhow::anyhow!(
"peer {} is not admitted for diagnostic 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 readiness_block.is_none() && 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 {
if let Some((_rejected, reason)) = readiness_block {
return Err(anyhow::anyhow!(
"peer {} did not become ready for diagnostic traffic: {}",
id,
reason
));
}
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<StateSnapshot> {
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<StateSnapshot> {
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 terminal = matches!(
snapshot.state.as_str(),
"closed" | "failed" | "disconnected"
);
let should_emit = {
let mut emitted = self
.last_emitted_connection_states
.lock()
.expect("wasm connection state cache poisoned");
if terminal {
emitted.insert(connection_id.to_string(), fingerprint);
true
} else {
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_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 Self::peer_session_is_stably_settled(snapshot_ref) {
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 async fn wait_for_settled_scope(
&self,
scope: &str,
timeout_ms: Option<u64>,
) -> Option<PeerSessionSnapshot> {
let scope = scope.trim();
if scope.is_empty() {
return None;
}
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 mut matches = self
.peer_sessions()
.await
.into_iter()
.filter(|snapshot| snapshot.scopes.iter().any(|candidate| candidate == scope));
let snapshot = matches.next();
let snapshot = if matches.next().is_none() {
snapshot
} else {
None
};
if let Some(snapshot_ref) = snapshot.as_ref() {
if Self::peer_session_is_stably_settled(snapshot_ref) {
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.filter(Self::peer_session_is_stably_settled);
}
#[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;
}
}
fn peer_session_is_stably_settled(snapshot: &PeerSessionSnapshot) -> bool {
if !snapshot.settled_ready {
return false;
}
if snapshot.candidate_connection_ids.is_empty()
&& snapshot.transport_generation <= 1
&& snapshot.route_generation <= 1
{
return true;
}
let stable_deadline_ms = snapshot
.last_lifecycle_transition_at_ms
.saturating_add(Self::SETTLED_PEER_STABILITY_WINDOW_MS as i64);
#[cfg(not(target_arch = "wasm32"))]
{
now_millis_i64() >= stable_deadline_ms
}
#[cfg(target_arch = "wasm32")]
{
js_sys::Date::now() as i64 >= stable_deadline_ms
}
}
pub(crate) fn peer_stream_generation_matches(
expected: &PeerSessionSnapshot,
current: &PeerSessionSnapshot,
opened_transport_stable_id: u64,
) -> bool {
Self::peer_session_is_stably_settled(current)
&& current.active_connection_id == expected.active_connection_id
&& current.node_id == expected.node_id
&& current.active_transport_stable_id == Some(opened_transport_stable_id)
&& current.transport_generation == expected.transport_generation
&& current.route_generation == expected.route_generation
}
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, send, recv) = self.open_peer_bi_current(id, timeout_ms).await?;
Ok((
snapshot.active_connection_id,
endpoint_id.to_string(),
send,
recv,
))
}
pub(crate) async fn open_peer_bi_current(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
PeerSessionSnapshot,
iroh::EndpointId,
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
let (snapshot, endpoint_id) = self.resolve_settled_peer_endpoint(id, timeout_ms).await?;
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
let (send, recv) = self
.open_generation_bound_peer_bi(
id,
&snapshot,
endpoint_id,
timeout_ms.map(std::time::Duration::from_millis),
"outgoing bidirectional application stream",
)
.await?;
Ok((snapshot, endpoint_id, send, recv))
}
async fn open_generation_bound_peer_bi(
&self,
id: &str,
expected: &PeerSessionSnapshot,
endpoint_id: iroh::EndpointId,
timeout: Option<std::time::Duration>,
operation: &str,
) -> anyhow::Result<(
crate::application_crypto_streams::PeerSendStream,
crate::application_crypto_streams::PeerRecvStream,
)> {
let expected_connection_id = expected.active_connection_id.as_deref().ok_or_else(|| {
anyhow::anyhow!("peer-stream-route-pending: peer {id} has no active connection")
})?;
let expected_transport_stable_id =
expected.active_transport_stable_id.ok_or_else(|| {
anyhow::anyhow!(
"peer-stream-route-pending: peer {id} has no active transport generation"
)
})?;
#[cfg(not(target_arch = "wasm32"))]
let opened = if let Some(timeout) = timeout {
tokio::time::timeout(
timeout,
self.open_current_bi_internal_with_transport_stable_id(endpoint_id),
)
.await
.map_err(|_| {
anyhow::anyhow!(
"peer-stream-route-pending: stream open timed out after {}ms for peer {id}",
timeout.as_millis(),
)
})?
} else {
self.open_current_bi_internal_with_transport_stable_id(endpoint_id)
.await
};
#[cfg(target_arch = "wasm32")]
let opened = {
let _ = timeout;
self.open_current_bi_internal_with_transport_stable_id(endpoint_id)
.await
};
let (opened_transport_stable_id, send, recv) = opened.map_err(|error| {
anyhow::anyhow!("peer-stream-route-pending: peer {id} has no current route: {error}")
})?;
if opened_transport_stable_id != expected_transport_stable_id {
anyhow::bail!(
"peer-stream-generation-stale: peer {id} physical generation changed while opening a protected stream"
);
}
let key = self
.required_crypto_key(Some(expected_connection_id), &endpoint_id, operation)
.await?;
let wrapped =
self.wrap_application_streams(expected_connection_id, key, send, recv, &[])?;
let current = self.peer_session(id).await;
let current_physical_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let generation_is_current = current.as_ref().is_some_and(|current| {
Self::peer_stream_generation_matches(expected, current, opened_transport_stable_id)
});
if current_physical_stable_id != Some(opened_transport_stable_id) || !generation_is_current
{
anyhow::bail!(
"peer-stream-generation-stale: peer {id} changed generation while opening a protected stream"
);
}
Ok(wrapped)
}
#[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))?;
let transport_stable_id = snapshot.active_transport_stable_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, transport_stable_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,
current.active_transport_stable_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_generation_bound_peer_bi(
id,
&snapshot,
endpoint_id,
Some(std::time::Duration::from_millis(
timeout_ms
.saturating_sub(
started
.elapsed()
.as_millis()
.try_into()
.unwrap_or(timeout_ms),
)
.max(1),
)),
"outgoing protected bidirectional application stream",
)
.await?;
Ok((
Some(connection_id.to_string()),
endpoint_id.to_string(),
send,
recv,
))
}
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_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 (mut send, recv) = self
.open_generation_bound_peer_bi(
id,
&snapshot,
endpoint_id,
timeout_ms.map(std::time::Duration::from_millis),
"outgoing explicit-file application stream",
)
.await?;
let connection_id = snapshot.active_connection_id.clone();
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_crypto_key(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
"outgoing diagnostic application stream",
)
.await?;
let connection_id = snapshot
.active_connection_id
.as_deref()
.ok_or_else(|| anyhow::anyhow!("diagnostic peer has no active connection"))?;
self.wrap_application_streams(connection_id, key, send, recv, &[])
.map(|wrapped| {
(
snapshot.active_connection_id,
remote_node_id,
wrapped.0,
wrapped.1,
)
})
}
pub async fn open_peer_bi_transport_only(
&self,
id: &str,
timeout_ms: Option<u64>,
) -> anyhow::Result<(
Option<String>,
String,
iroh::endpoint::SendStream,
iroh::endpoint::RecvStream,
)> {
let timeout_ms = timeout_ms.unwrap_or(10_000);
#[cfg(not(target_arch = "wasm32"))]
let started = std::time::Instant::now();
#[cfg(target_arch = "wasm32")]
let started_ms = js_sys::Date::now();
loop {
let snapshot = self
.peer_session(id)
.await
.ok_or_else(|| anyhow::anyhow!("peer {} not found", id))?;
let 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_crypto_key(
snapshot.active_connection_id.as_deref(),
&endpoint_id,
"outgoing unidirectional application stream",
)
.await?;
let send = match key {
Some(key) => {
let id = connection_id
.as_deref()
.ok_or_else(|| anyhow::anyhow!("peer has no active connection"))?;
let access = self.application_stream_access(id, key)?;
crate::application_crypto_streams::PeerSendStream::encrypted(send, key)
.with_access(access)
}
None => crate::application_crypto_streams::PeerSendStream::plain(send),
};
Ok((connection_id, remote_node_id, send))
}
pub const MAX_PEER_DATAGRAM_BYTES: usize = 1_024;
const PEER_DATAGRAM_MAGIC: [u8; 4] = *b"ORD1";
pub async fn send_peer_datagram(&self, id: &str, payload: &[u8]) -> anyhow::Result<()> {
if payload.len() > Self::MAX_PEER_DATAGRAM_BYTES {
anyhow::bail!(
"peer datagram payload exceeds {} bytes",
Self::MAX_PEER_DATAGRAM_BYTES
);
}
loop {
let (snapshot, endpoint_id) =
self.resolve_settled_peer_endpoint(id, Some(10_000)).await?;
let connection_id = snapshot
.active_connection_id
.as_deref()
.ok_or_else(|| anyhow::anyhow!("peer {id} has no active connection"))?;
let expected_stable_id = snapshot
.active_transport_stable_id
.ok_or_else(|| anyhow::anyhow!("peer {id} has no active transport generation"))?;
if self
.application_crypto_key_for_connection(Some(connection_id))
.is_none()
{
anyhow::bail!("protected peer datagrams require confirmed application crypto");
}
if self.connection_requires_application_crypto_confirmation(connection_id)
&& !self.connection_application_crypto_is_confirmed(
connection_id,
Some(expected_stable_id),
)
{
anyhow::bail!("protected peer datagrams require confirmed application crypto");
}
let protected = self.protect_outbound_application_payload(connection_id, payload)?;
let mut frame = Vec::with_capacity(Self::PEER_DATAGRAM_MAGIC.len() + protected.len());
frame.extend_from_slice(&Self::PEER_DATAGRAM_MAGIC);
frame.extend_from_slice(&protected);
let Some(connection) = self.get_connection(endpoint_id).await else {
if self
.wait_for_peer_transport_replacement(
id,
connection_id,
endpoint_id,
expected_stable_id,
10_000,
)
.await
.is_some()
{
continue;
}
anyhow::bail!("peer {id} has no live logical Iroh connection");
};
let stable_id = crate::transport_generation::for_connection(&connection);
if stable_id != expected_stable_id {
continue;
}
match connection
.send_datagram_wait(bytes::Bytes::from(frame))
.await
{
Ok(()) => return Ok(()),
Err(error) => {
if self
.wait_for_peer_transport_replacement(
id,
connection_id,
endpoint_id,
stable_id,
10_000,
)
.await
.is_some()
{
continue;
}
return Err(anyhow::anyhow!("peer datagram send failed: {error}"));
}
}
}
}
pub async fn send_peer_datagram_with_max_age(
&self,
id: &str,
payload: &[u8],
max_age_ms: u64,
) -> anyhow::Result<bool> {
if max_age_ms == 0 {
return Ok(false);
}
use futures::FutureExt;
let send = self.send_peer_datagram(id, payload).fuse();
#[cfg(target_arch = "wasm32")]
let deadline =
gloo_timers::future::sleep(std::time::Duration::from_millis(max_age_ms)).fuse();
#[cfg(not(target_arch = "wasm32"))]
let deadline = tokio::time::sleep(std::time::Duration::from_millis(max_age_ms)).fuse();
futures::pin_mut!(send, deadline);
futures::select! {
result = send => result.map(|()| true),
_ = deadline => Ok(false),
}
}
pub async fn receive_peer_datagram(&self, id: &str) -> anyhow::Result<Vec<u8>> {
loop {
let (snapshot, endpoint_id) =
self.resolve_settled_peer_endpoint(id, Some(10_000)).await?;
let connection_id = snapshot
.active_connection_id
.as_deref()
.ok_or_else(|| anyhow::anyhow!("peer {id} has no active connection"))?;
let expected_stable_id = snapshot
.active_transport_stable_id
.ok_or_else(|| anyhow::anyhow!("peer {id} has no active transport generation"))?;
if self
.application_crypto_key_for_connection(Some(connection_id))
.is_none()
|| (self.connection_requires_application_crypto_confirmation(connection_id)
&& !self.connection_application_crypto_is_confirmed(
connection_id,
Some(expected_stable_id),
))
{
anyhow::bail!("protected peer datagrams require confirmed application crypto");
}
let Some(connection) = self.get_connection(endpoint_id).await else {
if self
.wait_for_peer_transport_replacement(
id,
connection_id,
endpoint_id,
expected_stable_id,
10_000,
)
.await
.is_some()
{
continue;
}
anyhow::bail!("peer {id} has no live logical Iroh connection");
};
let stable_id = crate::transport_generation::for_connection(&connection);
if stable_id != expected_stable_id {
continue;
}
let frame = match connection.read_datagram().await {
Ok(frame) => frame,
Err(error) => {
let replacement = self
.wait_for_peer_transport_replacement(
id,
connection_id,
endpoint_id,
stable_id,
10_000,
)
.await
.is_some();
if replacement {
continue;
}
return Err(anyhow::anyhow!("peer datagram receive failed: {error}"));
}
};
let current_transport = self.get_connection(endpoint_id).await.filter(|current| {
crate::transport_generation::for_connection(current) == stable_id
});
let Ok((current_snapshot, current_endpoint_id)) =
self.resolve_settled_peer_endpoint(id, Some(0)).await
else {
continue;
};
if current_transport.is_none()
|| current_endpoint_id != endpoint_id
|| current_snapshot.active_connection_id.as_deref() != Some(connection_id)
|| current_snapshot.active_transport_stable_id != Some(stable_id)
{
continue;
}
let Some(protected) = frame.strip_prefix(&Self::PEER_DATAGRAM_MAGIC) else {
continue;
};
if let Ok(payload) = self.open_inbound_application_payload(connection_id, protected) {
return Ok(payload);
}
}
}
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());
}
return None;
}
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::ConnectionAdoption> {
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")
)
{
eprintln!(
"[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::ConnectionAdoption {
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"))]
{
eprintln!(
"[pluto-rtc][managed-retire] connection_id={} reason={:?}",
connection_id, normalized_reason
);
}
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
}
#[cfg(target_arch = "wasm32")]
pub(crate) 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;
};
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 recent_protocol_activity = super::auto_connect_impl::generation_activity_is_recent(
self.native_transport_protocol_activity_age(
&active_connection_id,
expected_transport_stable_id,
),
5_000,
);
let (health, healthy) = match probe {
_ if recent_protocol_activity => {
(crate::connection_manager::ConnectionHealth::Healthy, true)
}
Some(probe)
if probe.transport_stable_id == expected_transport_stable_id
&& (probe.responsive
|| super::auto_connect_impl::generation_activity_is_recent(
probe.last_inbound_activity_age,
5_000,
)) =>
{
(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_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;
}
eprintln!(
"[openrtc][presence][device-register] reason={} liveness_source=gateway-lease 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(()) => {
eprintln!(
"[openrtc][presence][device-register][ok] reason={} liveness_source=gateway-lease",
reason
);
true
}
Err(e) => {
eprintln!(
"[openrtc][presence][device-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(()) => {
eprintln!("[openrtc][presence][lease][ok] reason={}", reason);
true
}
Err(e) => {
eprintln!(
"[openrtc][presence][lease][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_ticket(&ticket);
let payload = suffix
.and_then(|value| crate::session_token::decode_payload(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_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_ticket(trimmed);
let Some(payload) =
suffix.and_then(|value| crate::session_token::decode_payload(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_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());
eprintln!(
"[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<()> {
let active_session_identity = self.active_session_identity();
#[cfg(not(target_arch = "wasm32"))]
let native_device_id = self
.native_device_identity
.read()
.await
.as_ref()
.map(|identity| identity.device_id.clone());
#[cfg(target_arch = "wasm32")]
let native_device_id: Option<String> = None;
let deleting_local_device = device_id_is_local(
device_id,
active_session_identity.as_ref(),
native_device_id.as_deref(),
);
if deleting_local_device {
self.stop_presence_loop();
}
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)
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn should_publish_excluded_peers_to_signaling(
external_coordination_active: bool,
) -> bool {
!external_coordination_active
}
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"))]
if !Self::should_publish_excluded_peers_to_signaling(
self.external_auto_connect_is_active().await,
) {
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"))]
if !Self::should_publish_excluded_peers_to_signaling(
self.external_auto_connect_is_active().await,
) {
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_local_manual_disconnect_status(snapshot);
snapshot = self.apply_device_status_admission_guard(snapshot);
self.enrich_device_status_latency(&mut snapshot).await;
snapshots.push(snapshot);
}
Ok(snapshots)
}
pub(crate) fn apply_local_manual_disconnect_status(
&self,
mut snapshot: DeviceStatusSnapshot,
) -> DeviceStatusSnapshot {
if !self.is_locally_auto_connect_excluded(&snapshot.device.device_id) {
return snapshot;
}
snapshot.connection_status = ConnectionStatus::Closed;
snapshot.settled_ready = false;
snapshot.readiness_state = ReadinessState::Closed;
snapshot.readiness_reason = crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string();
snapshot.peer_health = crate::connection_manager::ConnectionHealth::Stale;
snapshot.connection_id = None;
snapshot.active_transport_stable_id = None;
snapshot.parallel_transport = None;
snapshot.latency_ms = None;
snapshot.latency_by_transport = crate::client::LatencySnapshot::default();
snapshot
}
async fn enrich_device_status_latency(&self, snapshot: &mut DeviceStatusSnapshot) {
snapshot.latency_ms = None;
snapshot.latency_by_transport = crate::client::LatencySnapshot::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::LatencySnapshot::default();
if transport_labels.iter().flatten().any(|transport| {
crate::transport_label::normalize(transport).is_some_and(|label| {
crate::transport_label::is_iroh_base(label)
|| crate::transport_label::is_iroh_packet_carrier(label)
}) || 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);
}
}
}
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 = TransportStatus::from_config_with_ble_available(
&transport_config,
ble_available,
current_iroh_relay_provider(),
);
RuntimeStatus {
ready: true,
node_id,
transport,
}
}
pub fn product_capability_maturity(&self) -> ProductCapabilityMaturity {
ProductCapabilityMaturity::for_current_target()
}
#[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, None)
.await
}
pub async fn connect_known_device_with_token(
&self,
device_id: &str,
token: &str,
scope: &str,
max_connections: u32,
expires_at_ms: Option<u64>,
lookup_timeout_ms: Option<u64>,
) -> anyhow::Result<ManagedConnectResult> {
let device_id = device_id.trim();
let token = token.trim();
let scope = scope.trim();
if device_id.is_empty() {
return Err(anyhow::anyhow!(
"known-device admission requires a device id"
));
}
if token.is_empty() {
return Err(anyhow::anyhow!("known-device admission requires a token"));
}
if scope.is_empty() {
return Err(anyhow::anyhow!("known-device admission requires a scope"));
}
if expires_at_ms.is_some_and(|expiry| expiry <= crate::session_token::now_unix_ms()) {
return Err(anyhow::anyhow!(
"known-device admission expiry must be in the future"
));
}
let mut endpoint_revisions = self.known_device_endpoint_revision.subscribe();
let lookup = async {
loop {
if let Some(remote_node_id) = self.authoritative_node_for_device(device_id) {
let endpoint_id =
remote_node_id
.parse::<iroh::EndpointId>()
.map_err(|error| {
anyhow::anyhow!(
"known device {device_id} has an invalid authoritative endpoint: {error}"
)
})?;
if let Some(endpoint_addr) = self.cached_endpoint_addr(endpoint_id).await {
return Ok::<_, anyhow::Error>(endpoint_addr);
}
}
endpoint_revisions
.changed()
.await
.map_err(|_| anyhow::anyhow!("known-device endpoint directory closed"))?;
}
};
let lookup_timeout_ms = lookup_timeout_ms.unwrap_or(10_000).max(1);
#[cfg(not(target_arch = "wasm32"))]
let endpoint_addr = tokio::time::timeout(
std::time::Duration::from_millis(lookup_timeout_ms),
lookup,
)
.await
.map_err(|_| {
anyhow::anyhow!(
"timed out after {lookup_timeout_ms}ms waiting for Rust-owned endpoint state for known device {device_id}"
)
})??;
#[cfg(target_arch = "wasm32")]
let endpoint_addr = {
use futures::FutureExt;
let lookup = lookup.fuse();
let timeout =
gloo_timers::future::sleep(std::time::Duration::from_millis(lookup_timeout_ms))
.fuse();
futures::pin_mut!(lookup, timeout);
futures::select! {
endpoint_addr = lookup => endpoint_addr?,
_ = timeout => {
return Err(anyhow::anyhow!(
"timed out after {lookup_timeout_ms}ms waiting for Rust-owned endpoint state for known device {device_id}"
));
}
}
};
let remote_node_id = endpoint_addr.id.to_string();
let endpoint_ticket =
iroh_tickets::endpoint::EndpointTicket::new(endpoint_addr.clone()).to_string();
let compound_ticket = crate::session_token::expiring_ticket(
endpoint_ticket.as_str(),
token,
scope,
max_connections,
expires_at_ms,
);
let previous_repair_credential = self.cache_scoped_route_repair_credential(
&remote_node_id,
Some(crate::client::ScopedRouteRepairCredential {
token: token.to_string(),
scope: scope.to_string(),
max_connections,
expires_at_ms,
authoritative_device_id: device_id.to_string(),
}),
);
let first_connect = self.connect_device(None, compound_ticket.as_str()).await;
let mut first_error = None;
let mut result = match first_connect {
Ok(result) => Some(result),
Err(error) => {
first_error = Some(error);
None
}
};
if result
.as_ref()
.is_none_or(|result| result.approved_scope.as_deref() != Some(scope))
{
let stale_transport_stable_id = self
.get_connection(endpoint_addr.id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let retired_stale_transport = match stale_transport_stable_id {
Some(stable_id) => self
.disconnect_transport_generation_with_reason(
endpoint_addr.id,
stable_id,
crate::lifecycle_reason::REASON_SCOPED_ADMISSION_FRESH_REDIAL,
)
.await
.unwrap_or(false),
None => false,
};
eprintln!(
"[pluto-rtc][known-device-admission] retrying replacement scope on fresh transport device_id={} remote_node_id={} scope={} stale_transport_stable_id={:?} retired_stale_transport={} first_error={}",
device_id,
remote_node_id,
scope,
stale_transport_stable_id,
retired_stale_transport,
first_error
.as_ref()
.map(ToString::to_string)
.unwrap_or_else(|| "scope-not-approved".to_string()),
);
result = match self.connect_device(None, compound_ticket.as_str()).await {
Ok(result) => Some(result),
Err(error) => {
self.cache_scoped_route_repair_credential(
&remote_node_id,
previous_repair_credential,
);
return Err(match first_error {
Some(first_error) => anyhow::anyhow!(
"known-device admission for {device_id} failed before and after one fresh transport retry: first={first_error}; retry={error}"
),
None => error,
});
}
};
}
let mut result = result.expect("known-device admission retry produced a result");
if result.approved_scope.as_deref() != Some(scope) {
self.cache_scoped_route_repair_credential(&remote_node_id, previous_repair_credential);
return Err(anyhow::anyhow!(
"known-device admission for {device_id} did not receive approval for scope {scope}"
));
}
self.connection_manager
.set_device_id(&result.connection_id, device_id.to_string())
.await;
result.device_id = Some(device_id.to_string());
result.device_id_hint = Some(device_id.to_string());
Ok(result)
}
pub async fn observe_known_device_endpoint(
&self,
device_id: &str,
endpoint_ticket: &str,
) -> anyhow::Result<()> {
let device_id = device_id.trim();
if device_id.is_empty() {
return Err(anyhow::anyhow!(
"known-device endpoint observation requires a device id"
));
}
let endpoint_ticket = endpoint_ticket.trim();
let (iroh_ticket, _) = crate::session_token::split_ticket(endpoint_ticket);
let endpoint_addr = parse_endpoint_ticket(iroh_ticket)?;
let endpoint_id = endpoint_addr.id;
self.remember_endpoint_addr(&endpoint_addr).await;
self.reconcile_authoritative_device_node(device_id, &endpoint_id.to_string())
.await;
Ok(())
}
pub(crate) async fn connect_desired_device(
&self,
device_id: &str,
endpoint_ticket: &str,
desired_peer_evidence: &str,
) -> anyhow::Result<ManagedConnectResult> {
let (iroh_ticket, _) = crate::session_token::split_ticket(endpoint_ticket.trim());
let endpoint_addr = parse_endpoint_ticket(iroh_ticket)?;
let preferred_ticket =
self.preferred_scoped_route_ticket(endpoint_addr.id.to_string().as_str(), iroh_ticket);
self.connect_device_with_intent(
Some(device_id),
preferred_ticket.as_deref().unwrap_or(endpoint_ticket),
false,
Some(desired_peer_evidence),
)
.await
}
async fn connect_device_with_intent(
&self,
device_id: Option<&str>,
endpoint_ticket: &str,
clear_auto_connect_exclusion: bool,
desired_peer_evidence: Option<&str>,
) -> anyhow::Result<ManagedConnectResult> {
let (iroh_ticket, token_suffix) =
crate::session_token::split_ticket(endpoint_ticket.trim());
let extracted_token_payload = token_suffix
.and_then(|suffix| crate::session_token::decode_payload(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(not(target_arch = "wasm32"))]
let bilateral_user_device = extracted_token_payload
.as_ref()
.is_some_and(|payload| payload.scope.as_str() == "user-device");
#[cfg(not(target_arch = "wasm32"))]
let claimed_local_device_id = self
.auto_connect_loop_key
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.as_ref()
.map(|(_, device_id)| device_id.clone());
#[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);
#[cfg(not(target_arch = "wasm32"))]
if extracted_token_payload
.as_ref()
.is_some_and(|payload| payload.scope.as_str() == "user-device")
{
self.update_native_route_repair_credential(
&remote_node_id,
trimmed_device_id.as_deref(),
extracted_token.as_deref(),
extracted_token_suffix.as_deref(),
);
}
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"))]
eprintln!("{}", 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| crate::transport_generation::for_connection(&connection));
if let Some(live_transport_stable_id) = live_transport_stable_id {
self.commit_current_native_transport_record(
&connection_id,
&remote_node_id,
trimmed_device_id.clone(),
Some(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 current_remote_admission_scope = if let Some(token) = extracted_token.as_deref()
{
self.current_remote_session_admission_scope(
&connection_id,
endpoint_addr.id,
token,
)
.await
} else {
None
};
let has_remote_admission_proof = current_remote_admission_scope.is_some();
if let Some(scope) = current_remote_admission_scope {
approved_scope = Some(scope);
};
let needs_token_representation =
should_present_session_token_on_connected_transport(
extracted_token.is_some(),
has_remote_admission_proof,
);
#[cfg(not(target_arch = "wasm32"))]
let desired_presentation_decision = bilateral_user_device.then(|| {
let has_current_inbound_admission_proof =
live_transport_stable_id.is_some_and(|transport_stable_id| {
self.inbound_session_token_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
let peer_requested_reciprocal =
self.has_pending_reciprocal_session_admission_request(&connection_id);
super::auto_connect_impl::session_admission_presentation_decision(
extracted_token.is_some(),
has_remote_admission_proof,
has_current_inbound_admission_proof,
peer_requested_reciprocal,
clear_auto_connect_exclusion || local_node_id > remote_node_id,
)
});
#[cfg(not(target_arch = "wasm32"))]
let needs_token_representation = desired_presentation_decision
.map(|decision| decision.should_present)
.unwrap_or(needs_token_representation);
if needs_token_representation {
if let Some(token) = extracted_token.as_ref() {
eprintln!(
"[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())
);
#[cfg(not(target_arch = "wasm32"))]
let presentation_result =
if let Some(decision) = desired_presentation_decision {
self.present_and_accept_session_token_for_route_repair(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
claimed_local_device_id.clone(),
has_remote_admission_proof,
decision.request_reciprocal,
)
.await
} else {
self.present_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await
};
#[cfg(target_arch = "wasm32")]
let presentation_result = self
.present_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await;
match presentation_result {
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
eprintln!(
"[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 current = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
.ok_or_else(|| {
anyhow::anyhow!(
"managed connection was retired while existing transport settled"
)
})?;
if !matches!(
current.state,
crate::connection_manager::ConnectionState::Connected
) || current.transport_stable_id != live_transport_stable_id
|| !self.is_connection_transport_alive(endpoint_addr.id).await
{
return Err(anyhow::anyhow!(
"managed connection was superseded while existing transport settled"
));
}
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,
});
}
}
if clear_auto_connect_exclusion {
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
trimmed_device_id.clone(),
Some(remote_node_id.clone()),
)
.await;
self.connection_manager.set_connecting(&connection_id).await;
} else {
let automatic_device_id = trimmed_device_id.clone();
let current_user_device_assignment = if extracted_token_payload
.as_ref()
.is_some_and(|payload| payload.scope.as_str() == "user-device")
{
match (automatic_device_id.as_deref(), desired_peer_evidence) {
(Some(device_id), Some(evidence)) => {
self.current_desired_user_device_assignment_matches(
device_id,
&remote_node_id,
endpoint_ticket,
evidence,
)
.await
}
_ => false,
}
} else {
false
};
let began = self
.connection_manager
.begin_automatic_connect(
connection_id.clone(),
Some(remote_node_id.clone()),
trimmed_device_id.clone(),
Some(remote_node_id.clone()),
|requires_fresh_assignment| {
automatic_device_id
.as_deref()
.is_none_or(|device_id| !self.is_auto_connect_excluded(device_id))
&& (!requires_fresh_assignment
|| extracted_token_payload.as_ref().is_some_and(|payload| {
crate::session_token::matches_ticket(iroh_ticket, payload)
&& if payload.scope.as_str() == "user-device" {
current_user_device_assignment
} else {
self.session_token_registry.is_peer_assigned_to_scope(
payload.scope.as_str(),
&remote_node_id,
)
}
}))
},
)
.await;
if began.is_none() {
return Err(anyhow::anyhow!(
"automatic managed connect was cancelled by terminal peer intent"
));
}
if current_user_device_assignment {
let Some(device_id) = automatic_device_id.as_deref() else {
unreachable!("current user-device assignment requires a device id")
};
let Some(evidence) = desired_peer_evidence else {
unreachable!("current user-device assignment requires peer evidence")
};
if !self
.current_desired_user_device_assignment_matches(
device_id,
&remote_node_id,
endpoint_ticket,
evidence,
)
.await
{
return Err(anyhow::anyhow!(
"automatic managed connect was superseded by newer desired-peer evidence"
));
}
}
}
if let Some(device_id) = trimmed_device_id.as_deref() {
self.reconcile_authoritative_device_node(device_id, &remote_node_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| crate::transport_generation::for_connection(&connection));
log_phase("transport-record-start");
self.commit_current_native_transport_record(
&connection_id,
&remote_node_id,
trimmed_device_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 {
#[cfg(not(target_arch = "wasm32"))]
let desired_presentation_decision = bilateral_user_device.then(|| {
let has_remote_admission_proof =
transport_stable_id.is_some_and(|transport_stable_id| {
self.remote_session_token_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
let has_current_inbound_admission_proof =
transport_stable_id.is_some_and(|transport_stable_id| {
self.inbound_session_token_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
let peer_requested_reciprocal =
self.has_pending_reciprocal_session_admission_request(&connection_id);
super::auto_connect_impl::session_admission_presentation_decision(
true,
has_remote_admission_proof,
has_current_inbound_admission_proof,
peer_requested_reciprocal,
clear_auto_connect_exclusion || local_node_id > remote_node_id,
)
});
#[cfg(not(target_arch = "wasm32"))]
let should_present = desired_presentation_decision
.map(|decision| decision.should_present)
.unwrap_or(true);
#[cfg(target_arch = "wasm32")]
let should_present = true;
if should_present {
#[cfg(not(target_arch = "wasm32"))]
let presentation_result = if let Some(decision) = desired_presentation_decision {
self.present_and_accept_session_token_for_route_repair(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
claimed_local_device_id.clone(),
false,
decision.request_reciprocal,
)
.await
} else {
self.present_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await
};
#[cfg(target_arch = "wasm32")]
let presentation_result = self
.present_and_accept_session_token(
endpoint_addr.id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
trimmed_device_id.clone(),
)
.await;
match presentation_result {
Ok(approval_scope) => {
approved_scope = Some(approval_scope.clone());
eprintln!(
"[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
);
}
}
} else {
eprintln!(
"[pluto-rtc][connect-device][session-admission] awaiting elected presenter on desired transport connection_id={} remote_node_id={} local_node_id={}",
connection_id,
remote_node_id,
local_node_id,
);
}
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<ConnectionHealthView> {
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 => {
ManagedHealth::AwaitingReplacement
}
crate::connection_manager::ConnectionState::Connected => {
#[cfg(not(target_arch = "wasm32"))]
{
if admitted_settled_ready {
ManagedHealth::Healthy
} else if settle_deadline_elapsed {
ManagedHealth::Dead
} else {
ManagedHealth::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 =>
{
ManagedHealth::Healthy
}
_ if settle_deadline_elapsed => ManagedHealth::Dead,
crate::connection_manager::ConnectionHealth::Unknown
| crate::connection_manager::ConnectionHealth::Suspect
| crate::connection_manager::ConnectionHealth::Stale
| crate::connection_manager::ConnectionHealth::Healthy => {
ManagedHealth::AwaitingReplacement
}
}
}
}
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
| crate::connection_manager::ConnectionState::Failed => ManagedHealth::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 = ConnectionHealthView {
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 = ManagedHealth::Dead;
snapshot.readiness_state = ReadinessState::Failed;
} else {
snapshot.status = ManagedHealth::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<BridgeAction> {
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 {
ManagedHealth::Healthy => BridgeAction::Healthy,
ManagedHealth::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"
};
BridgeAction::Retire {
reason: timeout_reason.to_string(),
}
} else if snapshot
.node_id
.as_deref()
.and_then(|value| value.parse::<iroh::EndpointId>().ok())
.is_some()
{
BridgeAction::Rebind {
remote_node_id: snapshot.node_id,
transport_generation: snapshot.transport_generation,
}
} else {
BridgeAction::AwaitReplacement {
reason: snapshot.readiness_reason.clone(),
}
}
}
ManagedHealth::Dead => {
if settle_deadline_elapsed {
BridgeAction::Retire {
reason: "settle-timeout".to_string(),
}
} else {
BridgeAction::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_user_presence(
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 device and lease 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 =
crate::presence_policy::DEVICE_OFFLINE_RETENTION_MS as u64;
let mut durable_registered = false;
let mut live_registered = false;
let mut consecutive_publish_failures = 0_u32;
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(native_presence_retry_delay(consecutive_publish_failures));
retry_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) {
consecutive_publish_failures = 0;
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
);
}
} else {
let delay = native_presence_retry_delay(consecutive_publish_failures);
consecutive_publish_failures =
consecutive_publish_failures.saturating_add(1);
retry_interval.reset_after(delay);
}
}
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;
if !live_registered {
consecutive_publish_failures = 0;
retry_interval.reset_after(native_presence_retry_delay(0));
}
}
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;
if native_presence_is_ready(durable_registered, live_registered) {
consecutive_publish_failures = 0;
} else {
consecutive_publish_failures = 0;
retry_interval.reset_after(native_presence_retry_delay(0));
}
}
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] Rust browser presence is disabled; the TypeScript coordination gateway owns browser leases and device projection.",
));
}
Some(crate::presence::PresenceCommand::Stop) | None => {
break;
}
}
}
}
}
});
}
#[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"))]
eprintln!(
"[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();
self.scoped_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();
self.scoped_route_repair_credentials
.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 {} direct signaling path has been retired. Use the TypeScript coordination gateway adapter.",
operation
)
}
#[cfg(target_arch = "wasm32")]
fn warn_deprecated_browser_signaling(operation: &str) {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[OpenRTC] {} is deprecated in browsers. Use the TypeScript coordination gateway path.",
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"))]
eprintln!("[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"))]
eprintln!("[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; the TypeScript coordination gateway 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_auto_connect_wake.notify_waiters();
#[cfg(not(target_arch = "wasm32"))]
self.native_route_repair_credentials
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
self.scoped_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 _owner = match self.auto_connect_exclusion_owner.lock() {
Ok(owner) => owner,
Err(poisoned) => poisoned.into_inner(),
};
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 _owner = match self.auto_connect_exclusion_owner.lock() {
Ok(owner) => owner,
Err(poisoned) => poisoned.into_inner(),
};
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);
}
#[cfg(test)]
pub(crate) fn resume_peer_requested_auto_connect_if_authenticated(
&self,
device_id: &str,
) -> bool {
self.commit_peer_requested_reconnect_if_authenticated(device_id, || true)
}
pub(crate) fn commit_peer_requested_reconnect_if_authenticated(
&self,
device_id: &str,
commit_application_security: impl FnOnce() -> bool,
) -> bool {
let _owner = match self.auto_connect_exclusion_owner.lock() {
Ok(owner) => owner,
Err(poisoned) => poisoned.into_inner(),
};
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return false;
};
let mut excluded = match self.auto_connect_excluded.lock() {
Ok(excluded) => excluded,
Err(poisoned) => poisoned.into_inner(),
};
let mut aliases = match self.auto_connect_excluded_node_aliases.lock() {
Ok(aliases) => aliases,
Err(poisoned) => poisoned.into_inner(),
};
let mut peer_requested = match self.auto_connect_peer_requested_excluded.lock() {
Ok(peer_requested) => peer_requested,
Err(poisoned) => poisoned.into_inner(),
};
if !excluded.contains(&device_id) || !peer_requested.contains(&device_id) {
return false;
}
if !commit_application_security() {
return false;
}
peer_requested.remove(&device_id);
excluded.remove(&device_id);
aliases.remove(&device_id);
true
}
pub(crate) fn auto_connect_exclusion_for_peer(
&self,
device_id: Option<&str>,
node_id: Option<&str>,
) -> Option<bool> {
let _owner = match self.auto_connect_exclusion_owner.lock() {
Ok(owner) => owner,
Err(poisoned) => poisoned.into_inner(),
};
let device_id = device_id.and_then(Self::normalize_auto_connect_exclusion_key);
let node_id = node_id.and_then(Self::normalize_auto_connect_exclusion_key);
let excluded = match self.auto_connect_excluded.lock() {
Ok(excluded) => excluded,
Err(poisoned) => poisoned.into_inner(),
};
let aliases = match self.auto_connect_excluded_node_aliases.lock() {
Ok(aliases) => aliases,
Err(poisoned) => poisoned.into_inner(),
};
let peer_requested = match self.auto_connect_peer_requested_excluded.lock() {
Ok(peer_requested) => peer_requested,
Err(poisoned) => poisoned.into_inner(),
};
let excluded_device_id = device_id
.filter(|device_id| excluded.contains(device_id))
.or_else(|| {
node_id.as_deref().and_then(|node_id| {
aliases.iter().find_map(|(device_id, node_ids)| {
(excluded.contains(device_id) && node_ids.contains(node_id))
.then(|| device_id.clone())
})
})
})?;
Some(peer_requested.contains(&excluded_device_id))
}
pub fn is_auto_connect_excluded(&self, device_id: &str) -> bool {
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return false;
};
match self.auto_connect_excluded.lock() {
Ok(g) => g.contains(&device_id),
Err(p) => p.into_inner().contains(&device_id),
}
}
fn is_locally_auto_connect_excluded(&self, device_id: &str) -> bool {
let Some(device_id) = Self::normalize_auto_connect_exclusion_key(device_id) else {
return false;
};
let excluded = match self.auto_connect_excluded.lock() {
Ok(guard) => guard.contains(&device_id),
Err(poisoned) => poisoned.into_inner().contains(&device_id),
};
if !excluded {
return false;
}
match self.auto_connect_peer_requested_excluded.lock() {
Ok(guard) => !guard.contains(&device_id),
Err(poisoned) => !poisoned.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"))]
eprintln!(
"[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 {
eprintln!("[AppState] Returned to foreground — revisiting packet carriers 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
{
let capabilities = client
.native_peer_transport_capabilities
.read()
.await
.get(&record.connection_id)
.cloned()
.unwrap_or_default();
let _ = client
.maybe_start_preferred_native_iroh_carrier(
&record.connection_id,
Some(endpoint_id.as_str()),
capabilities.contains(
&crate::client::NativePeerTransportCapability::WebRtc,
),
capabilities.contains(
&crate::client::NativePeerTransportCapability::Moq,
),
)
.await;
}
}
});
}
}
}
pub fn is_app_backgrounded(&self) -> bool {
self.app_backgrounded
.load(std::sync::atomic::Ordering::Relaxed)
}
}
fn device_id_is_local(
device_id: &str,
active_session_identity: Option<&(String, String)>,
native_device_id: Option<&str>,
) -> bool {
active_session_identity.is_some_and(|(_, local_device_id)| local_device_id == device_id)
|| native_device_id.is_some_and(|local_device_id| local_device_id == device_id)
}
#[cfg(test)]
mod tests {
use super::{
device_id_is_local, native_presence_is_ready, native_presence_retry_delay,
native_presence_retry_targets, selected_iroh_latency_label, Client,
NativePresenceTicketPolicy,
};
#[tokio::test]
async fn pending_admission_cannot_project_terminal_connection_as_connecting() {
let client = Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
"test-app".to_string(),
Box::new(|| None),
);
let connection_id = "terminal-admission-projection";
client
.connection_manager
.upsert_pending(
connection_id.to_string(),
Some("terminal-node".to_string()),
Some("terminal-device".to_string()),
Some("terminal-node".to_string()),
)
.await;
client
.connection_manager
.commit_current_transport(
connection_id,
Some("terminal-node".to_string()),
71,
Some("initial".to_string()),
|_| {},
)
.await
.expect("connected transport");
client.set_connection_application_crypto_required(connection_id);
client
.connection_manager
.set_closed_if_current(
connection_id,
71,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await
.expect("terminal close");
let snapshot = client
.connection_state(connection_id)
.await
.expect("terminal snapshot");
assert_eq!(snapshot.state, "closed");
assert_eq!(snapshot.transport_state, "closed");
assert_eq!(snapshot.protocol_state, "closed");
assert_eq!(
snapshot.readiness_state,
crate::client::ReadinessState::Closed
);
assert!(!snapshot.routable);
}
#[test]
fn local_device_delete_detection_survives_missing_auto_connect_identity() {
let active = ("user".to_string(), "active-device".to_string());
assert!(device_id_is_local(
"active-device",
Some(&active),
Some("native-device"),
));
assert!(device_id_is_local(
"native-device",
None,
Some("native-device"),
));
assert!(!device_id_is_local(
"remote-device",
Some(&active),
Some("native-device"),
));
}
#[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));
}
#[test]
fn native_presence_retry_backoff_is_bounded() {
assert_eq!(
native_presence_retry_delay(0),
std::time::Duration::from_secs(2),
);
assert_eq!(
native_presence_retry_delay(1),
std::time::Duration::from_secs(4),
);
assert_eq!(
native_presence_retry_delay(5),
std::time::Duration::from_secs(60),
);
assert_eq!(
native_presence_retry_delay(100),
std::time::Duration::from_secs(60),
);
}
#[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_ticket(&ticket);
let payload = suffix
.and_then(|value| crate::session_token::decode_payload(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",
"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"),
);
assert_eq!(
selected_iroh_latency_label(
crate::client::IrohPathKind::WebRtc,
["iroh-relay", "webrtc"],
),
Some("webrtc"),
);
assert_eq!(
selected_iroh_latency_label(crate::client::IrohPathKind::Moq, ["moq", "iroh-relay"],),
Some("moq"),
);
}
}