use super::state_signaling_impl::should_present_session_token_on_connected_transport;
use super::*;
pub(crate) fn should_auto_present_session_token(
has_extracted_token: bool,
has_current_remote_admission_proof: bool,
) -> bool {
should_present_session_token_on_connected_transport(
has_extracted_token,
has_current_remote_admission_proof,
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SessionAdmissionPresentationDecision {
pub(crate) should_present: bool,
pub(crate) request_reciprocal: bool,
}
pub(crate) fn session_admission_presentation_decision(
has_extracted_token: bool,
has_current_remote_admission_proof: bool,
has_current_inbound_admission_proof: bool,
peer_requested_reciprocal: bool,
is_designated_initiator: bool,
) -> SessionAdmissionPresentationDecision {
if !has_extracted_token {
return SessionAdmissionPresentationDecision {
should_present: false,
request_reciprocal: false,
};
}
let bilateral_proof_is_current =
has_current_remote_admission_proof && has_current_inbound_admission_proof;
if bilateral_proof_is_current {
return SessionAdmissionPresentationDecision {
should_present: false,
request_reciprocal: false,
};
}
if !has_current_remote_admission_proof && !has_current_inbound_admission_proof {
return SessionAdmissionPresentationDecision {
should_present: is_designated_initiator,
request_reciprocal: is_designated_initiator,
};
}
if peer_requested_reciprocal {
return SessionAdmissionPresentationDecision {
should_present: true,
request_reciprocal: false,
};
}
let request_reciprocal = has_current_remote_admission_proof && is_designated_initiator;
SessionAdmissionPresentationDecision {
should_present: !has_current_remote_admission_proof || request_reciprocal,
request_reciprocal,
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ActiveConnectionProbeDecision {
Healthy,
PreserveLiveTransport,
Failed,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FailedProbeRetirementOutcome {
Retired,
DeferredAdmissionInFlight,
CancelledCurrentGenerationActivity,
GenerationChanged,
}
const ACTIVE_CONNECTION_PROBE_TIMEOUT_MS: i64 = 4_000;
pub(crate) fn generation_activity_is_recent(
last_inbound_activity_age: Option<std::time::Duration>,
probe_interval_ms: i64,
) -> bool {
let activity_window_ms = probe_interval_ms
.max(0)
.saturating_add(ACTIVE_CONNECTION_PROBE_TIMEOUT_MS) as u64;
last_inbound_activity_age
.is_some_and(|age| age <= std::time::Duration::from_millis(activity_window_ms))
}
#[cfg(not(target_arch = "wasm32"))]
impl Client {
pub(crate) fn record_native_transport_protocol_activity(
&self,
connection_id: &str,
transport_stable_id: u64,
) {
if let Ok(mut activity) = self.native_transport_protocol_activity.write() {
activity.insert(
connection_id.to_string(),
super::NativeTransportProtocolActivity {
transport_stable_id,
observed_at: std::time::Instant::now(),
},
);
}
}
pub(crate) fn native_transport_protocol_activity_age(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
) -> Option<std::time::Duration> {
let observed_at = self
.native_transport_protocol_activity
.read()
.ok()
.and_then(|activity| activity.get(connection_id).copied())
.filter(|activity| activity.transport_stable_id == expected_transport_stable_id)?
.observed_at;
Some(std::time::Instant::now().saturating_duration_since(observed_at))
}
pub(crate) fn forget_native_transport_protocol_activity_for_transport(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
) -> bool {
self.native_transport_protocol_activity
.write()
.ok()
.is_some_and(|mut activity| {
if activity
.get(connection_id)
.is_some_and(|entry| entry.transport_stable_id == expected_transport_stable_id)
{
activity.remove(connection_id);
true
} else {
false
}
})
}
pub(crate) fn forget_native_transport_protocol_activity(&self, connection_id: &str) {
if let Ok(mut activity) = self.native_transport_protocol_activity.write() {
activity.remove(connection_id);
}
}
async fn retire_transport_generation_after_failed_probe(
&self,
connection_id: &str,
endpoint_id: iroh::EndpointId,
expected_transport_stable_id: u64,
activity_window_ms: i64,
reason: &str,
) -> FailedProbeRetirementOutcome {
let Some(retirement_guard) = self
.session_token_registry
.try_begin_transport_retirement(connection_id)
else {
return FailedProbeRetirementOutcome::DeferredAdmissionInFlight;
};
let current_transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if current_transport_stable_id != Some(expected_transport_stable_id) {
drop(retirement_guard);
return FailedProbeRetirementOutcome::GenerationChanged;
}
if generation_activity_is_recent(
self.native_transport_protocol_activity_age(
connection_id,
expected_transport_stable_id,
),
activity_window_ms,
) {
drop(retirement_guard);
return FailedProbeRetirementOutcome::CancelledCurrentGenerationActivity;
}
let retired = self
.disconnect_transport_generation_with_reason(
endpoint_id,
expected_transport_stable_id,
reason,
)
.await
.unwrap_or(false);
drop(retirement_guard);
if retired {
FailedProbeRetirementOutcome::Retired
} else {
FailedProbeRetirementOutcome::GenerationChanged
}
}
}
#[cfg(target_arch = "wasm32")]
impl Client {
pub(crate) fn native_transport_protocol_activity_age(
&self,
_connection_id: &str,
_expected_transport_stable_id: u64,
) -> Option<std::time::Duration> {
None
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn active_connection_probe_decision(
active_probe_healthy: bool,
passive_transport_alive: bool,
) -> ActiveConnectionProbeDecision {
if active_probe_healthy {
ActiveConnectionProbeDecision::Healthy
} else if passive_transport_alive {
ActiveConnectionProbeDecision::PreserveLiveTransport
} else {
ActiveConnectionProbeDecision::Failed
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn active_connection_probe_decision_with_independent_route(
active_probe_healthy: bool,
passive_transport_alive: bool,
independent_route_ready: bool,
recent_generation_activity: bool,
) -> ActiveConnectionProbeDecision {
if independent_route_ready || recent_generation_activity {
ActiveConnectionProbeDecision::Healthy
} else {
active_connection_probe_decision(active_probe_healthy, passive_transport_alive)
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn next_active_connection_probe_failure_count(
decision: ActiveConnectionProbeDecision,
current_failures: u8,
) -> u8 {
if matches!(decision, ActiveConnectionProbeDecision::Healthy) {
0
} else {
current_failures.saturating_add(1).min(10)
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn active_connection_probe_is_healthy(
decision: ActiveConnectionProbeDecision,
failure_count: u8,
failure_threshold: u8,
) -> bool {
matches!(decision, ActiveConnectionProbeDecision::Healthy)
|| (matches!(
decision,
ActiveConnectionProbeDecision::PreserveLiveTransport
) && failure_count < failure_threshold)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SessionTokenPresentationFailureAction {
RetryOnNextTransport,
RejectAndDisconnect,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RetryableSessionAdmissionFailureDecision {
PreserveCurrentTransport,
RetireCurrentTransport,
ReplacementAlreadyWon,
}
pub(crate) fn superseded_desired_peer_was_withdrawn(
actor_is_current: bool,
current_revision: u64,
attempt_revision: u64,
peer_is_still_desired: bool,
) -> bool {
actor_is_current && current_revision > attempt_revision && !peer_is_still_desired
}
pub(crate) fn retryable_session_admission_failure_decision(
failure_count: u8,
failure_threshold: u8,
attempted_transport_stable_id: Option<u64>,
current_transport_stable_id: Option<u64>,
remote_runtime_instance_changed: bool,
) -> RetryableSessionAdmissionFailureDecision {
if attempted_transport_stable_id != current_transport_stable_id
|| attempted_transport_stable_id.is_none()
{
return RetryableSessionAdmissionFailureDecision::ReplacementAlreadyWon;
}
if !remote_runtime_instance_changed && failure_count < failure_threshold.max(1) {
RetryableSessionAdmissionFailureDecision::PreserveCurrentTransport
} else {
RetryableSessionAdmissionFailureDecision::RetireCurrentTransport
}
}
pub(crate) fn session_token_presentation_failure_action(
error: &str,
) -> SessionTokenPresentationFailureAction {
let error = error.to_ascii_lowercase();
if error.contains("connection lost")
|| error.contains("[session-token-response:timeout]")
|| error.contains("stream finished early")
|| error.contains("peer stream ended before")
|| error.contains("connection closed")
|| error.contains("stream reset")
|| error.contains("stream stopped")
|| error.contains("failed to open bi-stream")
|| error.contains("failed to write session-token presentation")
|| error.contains("missing connection for session-token presentation")
|| error.contains("[session-token-presentation:in-flight]")
|| error.contains("authoritative-device-binding-pending")
{
SessionTokenPresentationFailureAction::RetryOnNextTransport
} else {
SessionTokenPresentationFailureAction::RejectAndDisconnect
}
}
pub(crate) fn session_token_presentation_was_duplicate(error: &str) -> bool {
error
.to_ascii_lowercase()
.contains("[session-token-presentation:in-flight]")
}
fn remote_excludes_local_device(excluded_peers: &[String], local_device_id: &str) -> bool {
let local_device_id = local_device_id.trim();
!local_device_id.is_empty()
&& excluded_peers
.iter()
.any(|peer| peer.trim().eq_ignore_ascii_case(local_device_id))
}
fn observe_remote_runtime_instance(
last_known: &mut std::collections::HashMap<String, String>,
device_id: &str,
runtime_instance_id: Option<&str>,
) -> bool {
let Some(runtime_instance_id) = runtime_instance_id
.map(str::trim)
.filter(|value| !value.is_empty() && *value != device_id)
else {
return false;
};
let key = format!("{device_id}::runtime-instance");
last_known
.insert(key, runtime_instance_id.to_string())
.is_some_and(|previous| previous != runtime_instance_id)
}
#[cfg(not(target_arch = "wasm32"))]
fn runtime_instance_change_requires_replacement(
runtime_instance_changed: bool,
active_transport_responsive: bool,
) -> bool {
runtime_instance_changed && !active_transport_responsive
}
#[cfg(not(target_arch = "wasm32"))]
fn observe_remote_ticket(
last_known: &mut std::collections::HashMap<String, String>,
device_id: &str,
ticket: &str,
) -> bool {
let ticket = ticket.trim();
if ticket.is_empty() {
return false;
}
let key = format!("{device_id}::ticket");
let fingerprint = crate::session_token::token_fingerprint(ticket);
last_known
.insert(key, fingerprint.clone())
.is_some_and(|previous| previous != fingerprint)
}
#[cfg(not(target_arch = "wasm32"))]
fn reconcile_remote_ticket_retry_state(
last_known: &mut std::collections::HashMap<String, String>,
device_id: &str,
ticket: &str,
last_attempt_at: &mut std::collections::HashMap<String, i64>,
failure_count: &mut std::collections::HashMap<String, u8>,
failure_backoff_until: &mut std::collections::HashMap<String, i64>,
) -> bool {
if !observe_remote_ticket(last_known, device_id, ticket) {
return false;
}
last_attempt_at.remove(device_id);
clear_auto_connect_failure_state(device_id, failure_count, failure_backoff_until);
true
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct BrowserDesiredPeer {
pub(crate) device_id: String,
#[serde(default)]
pub(crate) node_id: Option<String>,
#[serde(default)]
pub(crate) ticket: Option<String>,
#[serde(default)]
pub(crate) online: bool,
#[serde(default)]
pub(crate) session_id: Option<String>,
#[serde(default)]
pub(crate) excluded_peers: Vec<String>,
}
#[derive(Debug, Default)]
struct BrowserDesiredPeerMailbox {
revision: u64,
peers: std::collections::BTreeMap<String, BrowserDesiredPeer>,
}
impl BrowserDesiredPeerMailbox {
fn accept(&mut self, revision: u64, peers: Vec<BrowserDesiredPeer>) -> bool {
self.accept_with_withdrawn(revision, peers).is_some()
}
fn accept_with_withdrawn(
&mut self,
revision: u64,
peers: Vec<BrowserDesiredPeer>,
) -> Option<Vec<BrowserDesiredPeer>> {
if revision <= self.revision {
return None;
}
let peers = peers
.into_iter()
.filter_map(|mut peer| {
peer.device_id = peer.device_id.trim().to_string();
if peer.device_id.is_empty() {
return None;
}
peer.node_id = peer
.node_id
.take()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
peer.ticket = peer
.ticket
.take()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
peer.session_id = peer
.session_id
.take()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
peer.excluded_peers = peer
.excluded_peers
.into_iter()
.map(|value| value.trim().to_ascii_lowercase())
.filter(|value| !value.is_empty())
.collect();
Some((peer.device_id.clone(), peer))
})
.collect::<std::collections::BTreeMap<_, _>>();
let withdrawn = self
.peers
.values()
.filter(|peer| {
peer.online
&& peers
.get(&peer.device_id)
.is_none_or(|current| !current.online)
})
.cloned()
.collect();
self.revision = revision;
self.peers = peers;
Some(withdrawn)
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Default)]
pub(crate) struct NativeExternalAutoConnectActorState {
active: bool,
generation: u64,
local_device_id: String,
mailbox: BrowserDesiredPeerMailbox,
peer_evidence: std::collections::HashMap<String, String>,
failure_count: std::collections::HashMap<String, u8>,
retry_not_before_ms: std::collections::HashMap<String, i64>,
non_initiator_not_before_ms: std::collections::HashMap<String, i64>,
scheduled_deadline_ms: Option<i64>,
schedule_epoch: u64,
reconcile_running: bool,
wake_pending: bool,
recovery_device_ids: std::collections::HashSet<String>,
recovery_suppressed_until_ms: std::collections::HashMap<String, i64>,
}
#[cfg(target_arch = "wasm32")]
struct BrowserAutoConnectActorState {
active: bool,
client: std::sync::Weak<Client>,
generation: u64,
local_device_id: String,
mailbox: BrowserDesiredPeerMailbox,
peer_evidence: std::collections::HashMap<String, String>,
failure_count: std::collections::HashMap<String, u8>,
retry_not_before_ms: std::collections::HashMap<String, i64>,
non_initiator_not_before_ms: std::collections::HashMap<String, i64>,
scheduled_deadline_ms: Option<i64>,
schedule_epoch: u64,
reconcile_running: bool,
wake_pending: bool,
}
#[cfg(target_arch = "wasm32")]
thread_local! {
static BROWSER_AUTO_CONNECT_ACTORS: std::cell::RefCell<
std::collections::HashMap<usize, std::rc::Rc<std::cell::RefCell<BrowserAutoConnectActorState>>>
> = std::cell::RefCell::new(std::collections::HashMap::new());
}
fn browser_auto_connect_failure_backoff_ms(failure_count: u8) -> i64 {
500_i64
.saturating_mul(1_i64 << (failure_count as u32).min(4))
.min(8_000)
}
fn browser_peer_evidence(peer: &BrowserDesiredPeer) -> String {
format!(
"{}|{}|{}",
peer.node_id.as_deref().unwrap_or_default(),
peer.ticket
.as_deref()
.map(crate::session_token::token_fingerprint)
.unwrap_or_default(),
peer.session_id.as_deref().unwrap_or_default(),
)
}
fn desired_peer_evidence_is_current(
mailbox: &BrowserDesiredPeerMailbox,
peer: &BrowserDesiredPeer,
) -> bool {
mailbox
.peers
.get(&peer.device_id)
.is_some_and(|current| browser_peer_evidence(current) == browser_peer_evidence(peer))
}
#[cfg(target_arch = "wasm32")]
fn browser_actor_key(client: &Arc<Client>) -> usize {
Arc::as_ptr(client) as usize
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BrowserDesiredPeerDecision {
ObserveExisting,
WaitUntil(i64),
Dial,
}
fn browser_desired_peer_decision(
local_node_id: &str,
remote_node_id: &str,
transport_alive: bool,
now_ms: i64,
non_initiator_not_before_ms: Option<i64>,
) -> BrowserDesiredPeerDecision {
if transport_alive {
return BrowserDesiredPeerDecision::ObserveExisting;
}
if local_node_id > remote_node_id {
return BrowserDesiredPeerDecision::Dial;
}
let not_before = non_initiator_not_before_ms.unwrap_or_else(|| now_ms.saturating_add(3_000));
if now_ms >= not_before {
BrowserDesiredPeerDecision::Dial
} else {
BrowserDesiredPeerDecision::WaitUntil(not_before)
}
}
fn should_replace_browser_retry_deadline(current: Option<i64>, next: Option<i64>) -> bool {
match (current, next) {
(None, None) => false,
(Some(_), None) | (None, Some(_)) => true,
(Some(current), Some(next)) => next < current,
}
}
fn browser_reconcile_deadlines(
outcome: BrowserPeerReconcileOutcome,
failure_retry_deadline: Option<i64>,
) -> (Option<i64>, Option<i64>) {
match outcome {
BrowserPeerReconcileOutcome::Retry => (failure_retry_deadline, None),
BrowserPeerReconcileOutcome::WakeAt(deadline) => (Some(deadline), None),
BrowserPeerReconcileOutcome::WaitUntil(deadline) => (None, Some(deadline)),
BrowserPeerReconcileOutcome::Stable | BrowserPeerReconcileOutcome::Ineligible => {
(None, None)
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BrowserPeerReconcileOutcome {
Stable,
Ineligible,
Retry,
WakeAt(i64),
WaitUntil(i64),
}
#[cfg(target_arch = "wasm32")]
fn browser_actor_is_current(actor: &BrowserAutoConnectActorState) -> bool {
actor.active
&& actor.client.upgrade().is_some_and(|client| {
client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
== actor.generation
})
}
#[cfg(target_arch = "wasm32")]
fn browser_attempt_is_current(
actor: &BrowserAutoConnectActorState,
_revision: u64,
peer: &BrowserDesiredPeer,
) -> bool {
browser_actor_is_current(actor) && desired_peer_evidence_is_current(&actor.mailbox, peer)
}
#[cfg(target_arch = "wasm32")]
fn deactivate_browser_actor(actor: &mut BrowserAutoConnectActorState) {
actor.active = false;
actor.mailbox.peers.clear();
actor.peer_evidence.clear();
actor.failure_count.clear();
actor.retry_not_before_ms.clear();
actor.non_initiator_not_before_ms.clear();
actor.scheduled_deadline_ms = None;
actor.schedule_epoch = actor.schedule_epoch.saturating_add(1);
actor.wake_pending = false;
}
#[cfg(target_arch = "wasm32")]
fn schedule_browser_auto_connect_retry(
actor: &std::rc::Rc<std::cell::RefCell<BrowserAutoConnectActorState>>,
) {
let (next_deadline, generation) = {
let state = actor.borrow();
if !browser_actor_is_current(&state) {
return;
}
let next = state
.retry_not_before_ms
.values()
.chain(state.non_initiator_not_before_ms.values())
.copied()
.min();
(next, state.generation)
};
let schedule = {
let mut state = actor.borrow_mut();
if !should_replace_browser_retry_deadline(state.scheduled_deadline_ms, next_deadline) {
return;
}
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = next_deadline;
next_deadline.map(|deadline| (deadline, state.schedule_epoch))
};
let Some((deadline, epoch)) = schedule else {
return;
};
let actor = actor.clone();
wasm_bindgen_futures::spawn_local(async move {
let delay_ms = deadline.saturating_sub(now_millis_i64()).max(0) as u64;
gloo_timers::future::sleep(std::time::Duration::from_millis(delay_ms)).await;
let should_wake = {
let mut state = actor.borrow_mut();
if !state.active
|| state.generation != generation
|| state.schedule_epoch != epoch
|| !browser_actor_is_current(&state)
{
false
} else {
state.scheduled_deadline_ms = None;
true
}
};
if should_wake {
wake_browser_auto_connect_actor(actor);
}
});
}
#[cfg(target_arch = "wasm32")]
fn wake_browser_auto_connect_actor(
actor: std::rc::Rc<std::cell::RefCell<BrowserAutoConnectActorState>>,
) {
let should_spawn = {
let mut state = actor.borrow_mut();
if !browser_actor_is_current(&state) {
return;
}
if state.reconcile_running {
state.wake_pending = true;
false
} else {
state.reconcile_running = true;
true
}
};
if !should_spawn {
return;
}
wasm_bindgen_futures::spawn_local(async move {
loop {
run_browser_auto_connect_pass(actor.clone()).await;
let run_again = {
let mut state = actor.borrow_mut();
if !browser_actor_is_current(&state) {
state.reconcile_running = false;
false
} else if state.wake_pending {
state.wake_pending = false;
true
} else {
state.reconcile_running = false;
false
}
};
if !run_again {
schedule_browser_auto_connect_retry(&actor);
break;
}
}
});
}
#[cfg(target_arch = "wasm32")]
async fn retire_stale_browser_attempt(
client: &Arc<Client>,
endpoint_id: iroh::EndpointId,
transport_stable_id: Option<u64>,
) {
if let Some(transport_stable_id) = transport_stable_id {
let node = client.iroh_node.read().await.as_ref().cloned();
if let Some(node) = node {
let _ = node
.disconnect_with_reason_if_current(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await;
}
}
}
#[cfg(target_arch = "wasm32")]
async fn reconcile_browser_desired_peer(
actor: &std::rc::Rc<std::cell::RefCell<BrowserAutoConnectActorState>>,
client: &Arc<Client>,
generation: u64,
revision: u64,
local_device_id: &str,
peer: &BrowserDesiredPeer,
non_initiator_not_before_ms: Option<i64>,
) -> BrowserPeerReconcileOutcome {
if client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= generation
{
return BrowserPeerReconcileOutcome::Ineligible;
}
if client.is_app_backgrounded() {
return BrowserPeerReconcileOutcome::Ineligible;
}
let parsed_ticket = peer.ticket.as_deref().and_then(|ticket| {
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(ticket);
parse_endpoint_ticket(iroh_ticket).ok()
});
let parsed_node_id = parsed_ticket.as_ref().map(|address| address.id.to_string());
let binding_node_id = peer.node_id.as_deref().or(parsed_node_id.as_deref());
if let Some(node_id) = binding_node_id {
client.bind_node_device_id(node_id, &peer.device_id).await;
}
if !actor
.try_borrow()
.is_ok_and(|state| browser_attempt_is_current(&state, revision, peer))
{
return BrowserPeerReconcileOutcome::Ineligible;
}
if !peer.online || peer.device_id == local_device_id {
return BrowserPeerReconcileOutcome::Ineligible;
}
if remote_excludes_local_device(&peer.excluded_peers, local_device_id) {
return BrowserPeerReconcileOutcome::Ineligible;
}
if client.is_auto_connect_excluded(&peer.device_id) {
return BrowserPeerReconcileOutcome::Ineligible;
}
let Some(ticket) = peer.ticket.as_deref() else {
return BrowserPeerReconcileOutcome::Ineligible;
};
let Some(endpoint_addr) = parsed_ticket else {
return BrowserPeerReconcileOutcome::Ineligible;
};
let endpoint_id = endpoint_addr.id;
let remote_node_id = endpoint_id.to_string();
if peer
.node_id
.as_deref()
.is_some_and(|node_id| node_id != remote_node_id)
{
return BrowserPeerReconcileOutcome::Ineligible;
}
let Some(local_node_id) = client.current_node_id().await else {
return BrowserPeerReconcileOutcome::WakeAt(now_millis_i64().saturating_add(250));
};
if local_node_id == remote_node_id {
return BrowserPeerReconcileOutcome::Ineligible;
}
client
.reconcile_authoritative_device_node(&peer.device_id, &remote_node_id)
.await;
if !actor
.try_borrow()
.is_ok_and(|state| browser_attempt_is_current(&state, revision, peer))
{
return BrowserPeerReconcileOutcome::Ineligible;
}
let transport_alive = client.is_connection_transport_alive(endpoint_id).await;
match browser_desired_peer_decision(
&local_node_id,
&remote_node_id,
transport_alive,
now_millis_i64(),
non_initiator_not_before_ms,
) {
BrowserDesiredPeerDecision::WaitUntil(deadline) => {
return BrowserPeerReconcileOutcome::WaitUntil(deadline);
}
BrowserDesiredPeerDecision::ObserveExisting | BrowserDesiredPeerDecision::Dial => {}
}
let result = client.connect_desired_device(&peer.device_id, ticket).await;
let transport_stable_id = client
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let (attempt_is_current, peer_was_withdrawn) = actor
.try_borrow()
.map(|state| {
let actor_is_current = browser_actor_is_current(&state);
(
actor_is_current && browser_attempt_is_current(&state, revision, peer),
superseded_desired_peer_was_withdrawn(
actor_is_current,
state.mailbox.revision,
revision,
state.mailbox.peers.contains_key(&peer.device_id),
),
)
})
.unwrap_or((false, false));
if !attempt_is_current {
if peer_was_withdrawn {
client
.retire_withdrawn_desired_peer_connections(&peer.device_id, peer.node_id.as_deref())
.await;
}
return BrowserPeerReconcileOutcome::Ineligible;
}
if client.is_auto_connect_excluded(&peer.device_id) {
retire_stale_browser_attempt(client, endpoint_id, transport_stable_id).await;
return BrowserPeerReconcileOutcome::Ineligible;
}
let Ok(result) = result else {
return BrowserPeerReconcileOutcome::Retry;
};
if matches!(result.state.as_str(), "failed" | "closed")
|| !client.is_connection_transport_alive(endpoint_id).await
{
return BrowserPeerReconcileOutcome::Retry;
}
let (iroh_ticket, token_suffix) = crate::session_token::split_compound_ticket(ticket);
if let Some(token) = token_suffix.and_then(|suffix| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, suffix)
.map(|payload| payload.token)
}) {
let has_proof = client
.has_current_remote_session_admission_proof(&result.connection_id, endpoint_id, &token)
.await;
if !has_proof {
return BrowserPeerReconcileOutcome::Retry;
}
}
BrowserPeerReconcileOutcome::Stable
}
#[cfg(target_arch = "wasm32")]
async fn run_browser_auto_connect_pass(
actor: std::rc::Rc<std::cell::RefCell<BrowserAutoConnectActorState>>,
) {
let (client, generation, revision, local_device_id, peers) = {
let state = actor.borrow();
if !browser_actor_is_current(&state) {
return;
}
let Some(client) = state.client.upgrade() else {
return;
};
(
client,
state.generation,
state.mailbox.revision,
state.local_device_id.clone(),
state.mailbox.peers.values().cloned().collect::<Vec<_>>(),
)
};
for peer in peers {
let (retry_not_before, non_initiator_not_before) = {
let state = actor.borrow();
if !browser_attempt_is_current(&state, revision, &peer) {
continue;
}
(
state.retry_not_before_ms.get(&peer.device_id).copied(),
state
.non_initiator_not_before_ms
.get(&peer.device_id)
.copied(),
)
};
let now = now_millis_i64();
if retry_not_before.is_some_and(|deadline| now < deadline) {
continue;
}
let outcome = reconcile_browser_desired_peer(
&actor,
&client,
generation,
revision,
&local_device_id,
&peer,
non_initiator_not_before,
)
.await;
let mut state = actor.borrow_mut();
if !browser_attempt_is_current(&state, revision, &peer) {
continue;
}
match outcome {
BrowserPeerReconcileOutcome::Stable => {
state.failure_count.remove(&peer.device_id);
state.retry_not_before_ms.remove(&peer.device_id);
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
BrowserPeerReconcileOutcome::Ineligible => {
state.failure_count.remove(&peer.device_id);
state.retry_not_before_ms.remove(&peer.device_id);
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
BrowserPeerReconcileOutcome::Retry => {
let count = state
.failure_count
.entry(peer.device_id.clone())
.or_insert(0);
*count = count.saturating_add(1).min(10);
let deadline = now_millis_i64()
.saturating_add(browser_auto_connect_failure_backoff_ms(*count));
let (retry_deadline, non_initiator_deadline) =
browser_reconcile_deadlines(outcome, Some(deadline));
match retry_deadline {
Some(deadline) => {
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
None => {
state.retry_not_before_ms.remove(&peer.device_id);
}
}
match non_initiator_deadline {
Some(deadline) => {
state
.non_initiator_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
None => {
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
}
}
BrowserPeerReconcileOutcome::WakeAt(deadline) => {
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
BrowserPeerReconcileOutcome::WaitUntil(deadline) => {
state.retry_not_before_ms.remove(&peer.device_id);
state
.non_initiator_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn native_external_attempt_is_current(
client: &Arc<Client>,
generation: u64,
_revision: u64,
peer: &BrowserDesiredPeer,
) -> bool {
if client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= generation
{
return false;
}
let state = client.external_desired_peer_actor.lock().await;
state.active
&& state.generation == generation
&& desired_peer_evidence_is_current(&state.mailbox, peer)
}
#[cfg(not(target_arch = "wasm32"))]
async fn retire_stale_native_external_attempt(
client: &Arc<Client>,
endpoint_id: iroh::EndpointId,
transport_stable_id: Option<u64>,
) {
if let Some(transport_stable_id) = transport_stable_id {
let node = client.iroh_node.read().await.as_ref().cloned();
if let Some(node) = node {
let _ = node
.disconnect_with_reason_if_current(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await;
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn reconcile_native_external_desired_peer(
client: &Arc<Client>,
generation: u64,
revision: u64,
local_device_id: &str,
peer: &BrowserDesiredPeer,
non_initiator_not_before_ms: Option<i64>,
) -> BrowserPeerReconcileOutcome {
if !native_external_attempt_is_current(client, generation, revision, peer).await
|| client.is_app_backgrounded()
{
return BrowserPeerReconcileOutcome::Ineligible;
}
let parsed_ticket = peer.ticket.as_deref().and_then(|ticket| {
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(ticket);
parse_endpoint_ticket(iroh_ticket).ok()
});
let parsed_node_id = parsed_ticket.as_ref().map(|address| address.id.to_string());
let binding_node_id = peer.node_id.as_deref().or(parsed_node_id.as_deref());
if let Some(node_id) = binding_node_id {
client.bind_node_device_id(node_id, &peer.device_id).await;
}
if !peer.online || peer.device_id == local_device_id {
return BrowserPeerReconcileOutcome::Ineligible;
}
if remote_excludes_local_device(&peer.excluded_peers, local_device_id) {
return BrowserPeerReconcileOutcome::Ineligible;
}
if client.is_auto_connect_excluded(&peer.device_id) {
return BrowserPeerReconcileOutcome::Ineligible;
}
let Some(ticket) = peer.ticket.as_deref() else {
return BrowserPeerReconcileOutcome::Ineligible;
};
let Some(endpoint_addr) = parsed_ticket else {
return BrowserPeerReconcileOutcome::Ineligible;
};
let endpoint_id = endpoint_addr.id;
let remote_node_id = endpoint_id.to_string();
if peer
.node_id
.as_deref()
.is_some_and(|node_id| node_id != remote_node_id)
{
return BrowserPeerReconcileOutcome::Ineligible;
}
let Some(local_node_id) = client.current_node_id().await else {
return BrowserPeerReconcileOutcome::WakeAt(now_millis_i64().saturating_add(250));
};
if local_node_id == remote_node_id {
return BrowserPeerReconcileOutcome::Ineligible;
}
client
.reconcile_authoritative_device_node(&peer.device_id, &remote_node_id)
.await;
if !native_external_attempt_is_current(client, generation, revision, peer).await {
return BrowserPeerReconcileOutcome::Ineligible;
}
let transport_alive = client.is_connection_transport_alive(endpoint_id).await;
match browser_desired_peer_decision(
&local_node_id,
&remote_node_id,
transport_alive,
now_millis_i64(),
non_initiator_not_before_ms,
) {
BrowserDesiredPeerDecision::WaitUntil(deadline) => {
return BrowserPeerReconcileOutcome::WaitUntil(deadline);
}
BrowserDesiredPeerDecision::ObserveExisting | BrowserDesiredPeerDecision::Dial => {}
}
let result = client.connect_desired_device(&peer.device_id, ticket).await;
let transport_stable_id = client
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if !native_external_attempt_is_current(client, generation, revision, peer).await {
let peer_was_withdrawn = {
let state = client.external_desired_peer_actor.lock().await;
let actor_is_current = state.active && state.generation == generation;
superseded_desired_peer_was_withdrawn(
actor_is_current,
state.mailbox.revision,
revision,
state.mailbox.peers.contains_key(&peer.device_id),
)
};
if peer_was_withdrawn {
client
.retire_withdrawn_desired_peer_connections(&peer.device_id, peer.node_id.as_deref())
.await;
}
return BrowserPeerReconcileOutcome::Ineligible;
}
if client.is_auto_connect_excluded(&peer.device_id) {
retire_stale_native_external_attempt(client, endpoint_id, transport_stable_id).await;
return BrowserPeerReconcileOutcome::Ineligible;
}
let Ok(result) = result else {
return BrowserPeerReconcileOutcome::Retry;
};
if matches!(result.state.as_str(), "failed" | "closed")
|| !client.is_connection_transport_alive(endpoint_id).await
{
return BrowserPeerReconcileOutcome::Retry;
}
let (iroh_ticket, token_suffix) = crate::session_token::split_compound_ticket(ticket);
if let Some(token) = token_suffix.and_then(|suffix| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, suffix)
.map(|payload| payload.token)
}) {
let has_proof = client
.has_current_remote_session_admission_proof(&result.connection_id, endpoint_id, &token)
.await;
if !has_proof {
return BrowserPeerReconcileOutcome::Retry;
}
}
BrowserPeerReconcileOutcome::Stable
}
#[cfg(not(target_arch = "wasm32"))]
async fn run_native_external_auto_connect_pass(client: Arc<Client>) {
let (generation, revision, local_device_id, peers, recovery_scoped) = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
{
return;
}
let recovery_device_ids = std::mem::take(&mut state.recovery_device_ids);
let recovery_scoped = !recovery_device_ids.is_empty();
let peers = state
.mailbox
.peers
.values()
.filter(|peer| {
recovery_device_ids.is_empty() || recovery_device_ids.contains(&peer.device_id)
})
.cloned()
.collect::<Vec<_>>();
(
state.generation,
state.mailbox.revision,
state.local_device_id.clone(),
peers,
recovery_scoped,
)
};
for peer in peers {
let (retry_not_before, non_initiator_not_before) = {
let state = client.external_desired_peer_actor.lock().await;
if !state.active
|| state.generation != generation
|| state.mailbox.revision != revision
|| !state
.mailbox
.peers
.get(&peer.device_id)
.is_some_and(|current| {
browser_peer_evidence(current) == browser_peer_evidence(&peer)
})
{
continue;
}
(
state.retry_not_before_ms.get(&peer.device_id).copied(),
state
.non_initiator_not_before_ms
.get(&peer.device_id)
.copied(),
)
};
let now = now_millis_i64();
if retry_not_before.is_some_and(|deadline| now < deadline) {
continue;
}
let outcome = reconcile_native_external_desired_peer(
&client,
generation,
revision,
&local_device_id,
&peer,
non_initiator_not_before,
)
.await;
if recovery_scoped {
println!(
"[pluto-rtc][auto-connect][recovery-pass] device_id={} outcome={:?}",
peer.device_id, outcome,
);
}
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| state.generation != generation
|| state.mailbox.revision != revision
|| !state
.mailbox
.peers
.get(&peer.device_id)
.is_some_and(|current| {
browser_peer_evidence(current) == browser_peer_evidence(&peer)
})
{
continue;
}
match outcome {
BrowserPeerReconcileOutcome::Stable | BrowserPeerReconcileOutcome::Ineligible => {
state.failure_count.remove(&peer.device_id);
state.retry_not_before_ms.remove(&peer.device_id);
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
BrowserPeerReconcileOutcome::Retry => {
let count = state
.failure_count
.entry(peer.device_id.clone())
.or_insert(0);
*count = count.saturating_add(1).min(10);
let deadline = now_millis_i64()
.saturating_add(browser_auto_connect_failure_backoff_ms(*count));
let (retry_deadline, non_initiator_deadline) =
browser_reconcile_deadlines(outcome, Some(deadline));
match retry_deadline {
Some(deadline) => {
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
None => {
state.retry_not_before_ms.remove(&peer.device_id);
}
}
match non_initiator_deadline {
Some(deadline) => {
state
.non_initiator_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
None => {
state.non_initiator_not_before_ms.remove(&peer.device_id);
}
}
if recovery_scoped {
state.recovery_device_ids.insert(peer.device_id.clone());
}
}
BrowserPeerReconcileOutcome::WakeAt(deadline) => {
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
state.non_initiator_not_before_ms.remove(&peer.device_id);
if recovery_scoped {
state.recovery_device_ids.insert(peer.device_id.clone());
}
}
BrowserPeerReconcileOutcome::WaitUntil(deadline) => {
state.retry_not_before_ms.remove(&peer.device_id);
state
.non_initiator_not_before_ms
.insert(peer.device_id.clone(), deadline);
if recovery_scoped {
state.recovery_device_ids.insert(peer.device_id.clone());
}
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn schedule_native_external_auto_connect_retry(client: Arc<Client>) {
let schedule = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
{
return;
}
let next = state
.retry_not_before_ms
.values()
.chain(state.non_initiator_not_before_ms.values())
.copied()
.min();
if !should_replace_browser_retry_deadline(state.scheduled_deadline_ms, next) {
return;
}
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = next;
next.map(|deadline| (deadline, state.schedule_epoch, state.generation))
};
let Some((deadline, epoch, generation)) = schedule else {
return;
};
tokio::spawn(async move {
let delay_ms = deadline.saturating_sub(now_millis_i64()).max(0) as u64;
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
let should_wake = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| state.generation != generation
|| state.schedule_epoch != epoch
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= generation
{
false
} else {
state.scheduled_deadline_ms = None;
true
}
};
if should_wake {
let _ = wake_native_external_auto_connect_boxed(client).await;
}
});
}
#[cfg(not(target_arch = "wasm32"))]
fn wake_native_external_auto_connect_boxed(
client: Arc<Client>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send>> {
Box::pin(async move {
let should_spawn = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
{
return false;
}
if state.reconcile_running {
state.wake_pending = true;
false
} else {
state.reconcile_running = true;
true
}
};
if !should_spawn {
return true;
}
tokio::spawn(async move {
loop {
run_native_external_auto_connect_pass(client.clone()).await;
let run_again = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
{
state.reconcile_running = false;
false
} else if state.wake_pending {
state.wake_pending = false;
true
} else {
state.reconcile_running = false;
false
}
};
if !run_again {
schedule_native_external_auto_connect_retry(client.clone()).await;
break;
}
}
});
true
})
}
#[cfg(not(target_arch = "wasm32"))]
impl Client {
#[doc(hidden)]
pub async fn external_desired_peer_debug_snapshot(&self) -> serde_json::Value {
let state = self.external_desired_peer_actor.lock().await;
let peers = state
.mailbox
.peers
.values()
.map(|peer| {
let ticket_node_id = peer.ticket.as_deref().and_then(|ticket| {
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(ticket);
parse_endpoint_ticket(iroh_ticket)
.ok()
.map(|address| address.id.to_string())
});
serde_json::json!({
"deviceId": peer.device_id,
"nodeId": peer.node_id,
"ticketNodeId": ticket_node_id,
"ticketFingerprint": peer.ticket.as_deref()
.map(crate::session_token::token_fingerprint),
"online": peer.online,
"sessionId": peer.session_id,
})
})
.collect::<Vec<_>>();
serde_json::json!({
"active": state.active,
"generation": state.generation,
"revision": state.mailbox.revision,
"reconcileRunning": state.reconcile_running,
"wakePending": state.wake_pending,
"scheduledDeadlineMs": state.scheduled_deadline_ms,
"lastFailureCounts": state.failure_count,
"retryNotBeforeMs": state.retry_not_before_ms,
"peers": peers,
})
}
pub async fn start_external_auto_connect(
self: &Arc<Self>,
user_id: String,
local_device_id: String,
) -> anyhow::Result<()> {
let user_id = user_id.trim().to_string();
let local_device_id = local_device_id.trim().to_string();
if user_id.is_empty() || local_device_id.is_empty() {
return Err(anyhow::anyhow!(
"external auto-connect requires non-empty user and local device ids"
));
}
let key = (user_id.clone(), local_device_id.clone());
let duplicate_start = {
let state = self.external_desired_peer_actor.lock().await;
let guard = self
.auto_connect_loop_key
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
guard.as_ref() == Some(&key)
&& state.active
&& state.local_device_id == local_device_id
&& self
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
== state.generation
};
if duplicate_start {
self.wake_native_external_auto_connect().await;
return Ok(());
}
let generation = {
let mut guard = self
.auto_connect_loop_key
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
self.known_device_ids_by_node
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
self.native_route_repair_credentials
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
*guard = Some(key);
self.auto_connect_generation
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1
};
let mut state = self.external_desired_peer_actor.lock().await;
*state = NativeExternalAutoConnectActorState {
active: true,
generation,
local_device_id,
..NativeExternalAutoConnectActorState::default()
};
Ok(())
}
pub async fn submit_external_desired_peers(
self: &Arc<Self>,
revision: u64,
peers_json: &str,
) -> anyhow::Result<bool> {
if peers_json.len() > 256 * 1024 {
return Err(anyhow::anyhow!(
"external desired-peer payload exceeds 256 KiB"
));
}
let peers: Vec<BrowserDesiredPeer> = serde_json::from_str(peers_json)
.map_err(|error| anyhow::anyhow!("invalid external desired-peer payload: {error}"))?;
if peers.len() > 100 {
return Err(anyhow::anyhow!(
"external desired-peer payload exceeds 100 peers"
));
}
let (bindings, withdrawn) = {
let mut state = self.external_desired_peer_actor.lock().await;
if !state.active {
return Err(anyhow::anyhow!("external auto-connect is not started"));
}
let Some(withdrawn) = state.mailbox.accept_with_withdrawn(revision, peers) else {
return Ok(false);
};
let desired_ids = state
.mailbox
.peers
.keys()
.cloned()
.collect::<std::collections::HashSet<_>>();
state.peer_evidence.retain(|id, _| desired_ids.contains(id));
state.failure_count.retain(|id, _| desired_ids.contains(id));
state
.retry_not_before_ms
.retain(|id, _| desired_ids.contains(id));
state
.non_initiator_not_before_ms
.retain(|id, _| desired_ids.contains(id));
let evidence = state
.mailbox
.peers
.iter()
.map(|(id, peer)| (id.clone(), browser_peer_evidence(peer)))
.collect::<Vec<_>>();
for (id, evidence) in evidence {
let changed = state
.peer_evidence
.insert(id.clone(), evidence.clone())
.is_some_and(|previous| previous != evidence);
if changed {
state.failure_count.remove(&id);
state.retry_not_before_ms.remove(&id);
state.non_initiator_not_before_ms.remove(&id);
}
}
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = None;
let bindings = state
.mailbox
.peers
.values()
.filter(|peer| peer.online)
.filter_map(|peer| {
let ticket_node_id = peer.ticket.as_deref().and_then(|ticket| {
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(ticket);
parse_endpoint_ticket(iroh_ticket)
.ok()
.map(|address| address.id.to_string())
});
let node_id = match (peer.node_id.as_deref(), ticket_node_id.as_deref()) {
(Some(directory), Some(ticket)) if directory != ticket => return None,
(Some(directory), _) => directory.to_string(),
(None, Some(ticket)) => ticket.to_string(),
(None, None) => return None,
};
Some((peer.device_id.clone(), node_id))
})
.collect::<Vec<_>>();
(bindings, withdrawn)
};
for peer in withdrawn {
let still_withdrawn = {
let state = self.external_desired_peer_actor.lock().await;
state.active
&& state.mailbox.revision == revision
&& state
.mailbox
.peers
.get(&peer.device_id)
.is_none_or(|current| !current.online)
};
if still_withdrawn {
self.retire_withdrawn_desired_peer_connections(
&peer.device_id,
peer.node_id.as_deref(),
)
.await;
}
}
for (device_id, node_id) in bindings {
self.observe_authoritative_device_node(&device_id, &node_id);
}
self.wake_native_external_auto_connect().await;
Ok(true)
}
pub async fn wake_native_external_auto_connect(&self) -> bool {
{
let mut state = self.external_desired_peer_actor.lock().await;
state.recovery_device_ids.clear();
state.recovery_suppressed_until_ms.clear();
}
wake_native_external_auto_connect_boxed(Arc::new(self.clone())).await
}
pub(crate) async fn wake_native_external_auto_connect_for_replaced_route(
&self,
remote_node_id: &str,
) -> bool {
let remote_node_id = remote_node_id.trim();
let Some(remote_device_id) = self
.known_device_ids_by_node
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(remote_node_id)
.cloned()
else {
println!(
"[pluto-rtc][auto-connect][replacement-route] remote_node_id={} action=skip reason=device-binding-missing",
remote_node_id,
);
return false;
};
{
let mut state = self.external_desired_peer_actor.lock().await;
if !state.active
|| self
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
|| !state.mailbox.peers.contains_key(&remote_device_id)
{
println!(
"[pluto-rtc][auto-connect][replacement-route] remote_node_id={} device_id={} action=skip reason=actor-or-peer-inactive",
remote_node_id, remote_device_id,
);
return false;
}
state.recovery_device_ids.insert(remote_device_id.clone());
state.recovery_suppressed_until_ms.remove(&remote_device_id);
state.retry_not_before_ms.remove(&remote_device_id);
state.non_initiator_not_before_ms.remove(&remote_device_id);
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = None;
}
println!(
"[pluto-rtc][auto-connect][replacement-route] remote_node_id={} device_id={} action=wake",
remote_node_id, remote_device_id,
);
wake_native_external_auto_connect_boxed(Arc::new(self.clone())).await
}
pub async fn schedule_native_external_auto_connect_recovery(
&self,
remote_node_id: &str,
) -> bool {
const RECOVERY_COALESCE_MS: i64 = 1_000;
const RECOVERY_DUPLICATE_SUPPRESSION_MS: i64 = MANAGED_SETTLE_DEADLINE_MS;
let remote_node_id = remote_node_id.trim();
let Some(remote_device_id) = self
.known_device_ids_by_node
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(remote_node_id)
.cloned()
else {
return false;
};
let schedule = {
let mut state = self.external_desired_peer_actor.lock().await;
if !state.active
|| self
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= state.generation
{
return false;
}
let now = now_millis_i64();
if state
.recovery_suppressed_until_ms
.get(&remote_device_id)
.is_some_and(|deadline| *deadline > now)
{
return true;
}
state.recovery_device_ids.insert(remote_device_id.clone());
state.recovery_suppressed_until_ms.insert(
remote_device_id,
now.saturating_add(RECOVERY_DUPLICATE_SUPPRESSION_MS),
);
let deadline = now.saturating_add(RECOVERY_COALESCE_MS);
if state
.scheduled_deadline_ms
.is_some_and(|scheduled| scheduled <= deadline)
{
return true;
}
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = Some(deadline);
(deadline, state.schedule_epoch, state.generation)
};
let client = Arc::new(self.clone());
tokio::spawn(async move {
let (deadline, epoch, generation) = schedule;
let delay_ms = deadline.saturating_sub(now_millis_i64()).max(0) as u64;
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
let should_wake = {
let mut state = client.external_desired_peer_actor.lock().await;
if !state.active
|| state.generation != generation
|| state.schedule_epoch != epoch
|| client
.auto_connect_generation
.load(std::sync::atomic::Ordering::SeqCst)
!= generation
{
false
} else {
state.scheduled_deadline_ms = None;
true
}
};
if should_wake {
let _ = wake_native_external_auto_connect_boxed(client).await;
}
});
true
}
pub(crate) async fn external_auto_connect_is_active(&self) -> bool {
self.external_desired_peer_actor.lock().await.active
}
pub async fn stop_external_auto_connect(&self) {
let mut state = self.external_desired_peer_actor.lock().await;
state.active = false;
state.mailbox.peers.clear();
state.peer_evidence.clear();
state.failure_count.clear();
state.retry_not_before_ms.clear();
state.non_initiator_not_before_ms.clear();
state.scheduled_deadline_ms = None;
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.wake_pending = false;
state.recovery_device_ids.clear();
state.recovery_suppressed_until_ms.clear();
}
}
#[cfg(target_arch = "wasm32")]
impl Client {
pub(crate) fn start_browser_auto_connect(
self: &Arc<Self>,
user_id: String,
local_device_id: String,
) -> anyhow::Result<()> {
let user_id = user_id.trim().to_string();
let local_device_id = local_device_id.trim().to_string();
if user_id.is_empty() || local_device_id.is_empty() {
return Err(anyhow::anyhow!(
"browser auto-connect requires non-empty user and local device ids"
));
}
let key = (user_id, local_device_id.clone());
let actor_key = browser_actor_key(self);
let duplicate_actor = BROWSER_AUTO_CONNECT_ACTORS.with(|actors| {
actors.borrow().get(&actor_key).cloned().filter(|actor| {
let state = actor.borrow();
browser_actor_is_current(&state) && state.local_device_id == local_device_id
})
});
{
let guard = self
.auto_connect_loop_key
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if guard.as_ref() == Some(&key) {
if let Some(actor) = duplicate_actor {
wake_browser_auto_connect_actor(actor);
return Ok(());
}
}
}
let generation = {
let mut guard = self
.auto_connect_loop_key
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
self.known_device_ids_by_node
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clear();
*guard = Some(key);
self.auto_connect_generation
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1
};
let actor = std::rc::Rc::new(std::cell::RefCell::new(BrowserAutoConnectActorState {
active: true,
client: Arc::downgrade(self),
generation,
local_device_id,
mailbox: BrowserDesiredPeerMailbox::default(),
peer_evidence: std::collections::HashMap::new(),
failure_count: std::collections::HashMap::new(),
retry_not_before_ms: std::collections::HashMap::new(),
non_initiator_not_before_ms: std::collections::HashMap::new(),
scheduled_deadline_ms: None,
schedule_epoch: 0,
reconcile_running: false,
wake_pending: false,
}));
BROWSER_AUTO_CONNECT_ACTORS.with(|actors| {
if let Some(previous) = actors.borrow_mut().insert(actor_key, actor) {
deactivate_browser_actor(&mut previous.borrow_mut());
}
});
Ok(())
}
pub(crate) fn submit_browser_desired_peers(
self: &Arc<Self>,
revision: u64,
peers_json: &str,
) -> anyhow::Result<bool> {
let peers: Vec<BrowserDesiredPeer> = serde_json::from_str(peers_json)
.map_err(|error| anyhow::anyhow!("invalid browser desired-peer payload: {error}"))?;
let actor = BROWSER_AUTO_CONNECT_ACTORS
.with(|actors| actors.borrow().get(&browser_actor_key(self)).cloned());
let Some(actor) = actor else {
return Err(anyhow::anyhow!("browser auto-connect is not started"));
};
let (accepted, current_bindings, withdrawn) = {
let mut state = actor.borrow_mut();
if !browser_actor_is_current(&state) {
(false, Vec::new(), Vec::new())
} else if let Some(withdrawn) = state.mailbox.accept_with_withdrawn(revision, peers) {
let desired_ids = state
.mailbox
.peers
.keys()
.cloned()
.collect::<std::collections::HashSet<_>>();
state
.peer_evidence
.retain(|device_id, _| desired_ids.contains(device_id));
state
.failure_count
.retain(|device_id, _| desired_ids.contains(device_id));
state
.retry_not_before_ms
.retain(|device_id, _| desired_ids.contains(device_id));
state
.non_initiator_not_before_ms
.retain(|device_id, _| desired_ids.contains(device_id));
let evidence = state
.mailbox
.peers
.iter()
.map(|(device_id, peer)| (device_id.clone(), browser_peer_evidence(peer)))
.collect::<Vec<_>>();
for (device_id, evidence) in evidence {
let changed = state
.peer_evidence
.insert(device_id.clone(), evidence.clone())
.is_some_and(|previous| previous != evidence);
if changed {
state.failure_count.remove(&device_id);
state.retry_not_before_ms.remove(&device_id);
state.non_initiator_not_before_ms.remove(&device_id);
}
}
state.schedule_epoch = state.schedule_epoch.saturating_add(1);
state.scheduled_deadline_ms = None;
let bindings = state
.mailbox
.peers
.values()
.filter(|peer| peer.online)
.filter_map(|peer| {
let ticket_node_id = peer.ticket.as_deref().and_then(|ticket| {
let (iroh_ticket, _) =
crate::session_token::split_compound_ticket(ticket);
parse_endpoint_ticket(iroh_ticket)
.ok()
.map(|address| address.id.to_string())
});
let node_id = match (peer.node_id.as_deref(), ticket_node_id.as_deref()) {
(Some(directory_node_id), Some(ticket_node_id))
if directory_node_id != ticket_node_id =>
{
return None;
}
(Some(directory_node_id), _) => directory_node_id.to_string(),
(None, Some(ticket_node_id)) => ticket_node_id.to_string(),
(None, None) => return None,
};
Some((peer.device_id.clone(), node_id))
})
.collect::<Vec<_>>();
(true, bindings, withdrawn)
} else {
(false, Vec::new(), Vec::new())
}
};
if accepted {
for (device_id, node_id) in current_bindings {
self.observe_authoritative_device_node(&device_id, &node_id);
}
if !withdrawn.is_empty() {
let client = Arc::clone(self);
let actor_for_retirement = actor.clone();
wasm_bindgen_futures::spawn_local(async move {
for peer in withdrawn {
let still_withdrawn = {
let state = actor_for_retirement.borrow();
browser_actor_is_current(&state)
&& state.mailbox.revision == revision
&& state
.mailbox
.peers
.get(&peer.device_id)
.is_none_or(|current| !current.online)
};
if still_withdrawn {
client
.retire_withdrawn_desired_peer_connections(
&peer.device_id,
peer.node_id.as_deref(),
)
.await;
}
}
});
}
wake_browser_auto_connect_actor(actor);
}
Ok(accepted)
}
pub(crate) fn wake_browser_auto_connect(self: &Arc<Self>) -> bool {
let actor = BROWSER_AUTO_CONNECT_ACTORS
.with(|actors| actors.borrow().get(&browser_actor_key(self)).cloned());
let Some(actor) = actor else {
return false;
};
if !browser_actor_is_current(&actor.borrow()) {
return false;
}
wake_browser_auto_connect_actor(actor);
true
}
pub(crate) fn stop_browser_auto_connect(self: &Arc<Self>) {
let actor = BROWSER_AUTO_CONNECT_ACTORS
.with(|actors| actors.borrow_mut().remove(&browser_actor_key(self)));
if let Some(actor) = actor {
deactivate_browser_actor(&mut actor.borrow_mut());
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn replacement_snapshot(
last_reconnect_reason: Option<&str>,
) -> ManagedConnectionHealthSnapshot {
ManagedConnectionHealthSnapshot {
connection_id: "conn-1".to_string(),
device_id: Some("device-1".to_string()),
device_id_hint: Some("device-1".to_string()),
node_id: Some("node-1".to_string()),
active_transport_stable_id: None,
transport_generation: 2,
route_generation: 0,
status: ManagedConnectionHealthStatus::AwaitingReplacement,
settled_ready: false,
readiness_state: ReadinessState::AwaitingReplacement,
replacement_pending: true,
last_lifecycle_transition_at_ms: 1_000,
readiness_reason: "missing-active-transport".to_string(),
transition_count: 1,
connecting_transition_count: 0,
replacement_count: 1,
retire_count: 0,
last_disconnect_reason: None,
last_reconnect_reason: last_reconnect_reason.map(str::to_string),
}
}
#[test]
fn runtime_instance_change_is_observed_only_after_the_initial_generation() {
let mut last_known = std::collections::HashMap::new();
assert!(!observe_remote_runtime_instance(
&mut last_known,
"device-1",
Some("runtime-1"),
));
assert!(!observe_remote_runtime_instance(
&mut last_known,
"device-1",
Some("runtime-1"),
));
assert!(observe_remote_runtime_instance(
&mut last_known,
"device-1",
Some("runtime-2"),
));
}
#[cfg(not(target_arch = "wasm32"))]
#[test]
fn delayed_runtime_instance_update_preserves_a_responsive_replacement() {
assert!(!runtime_instance_change_requires_replacement(true, true));
}
#[cfg(not(target_arch = "wasm32"))]
#[test]
fn runtime_instance_update_replaces_an_unresponsive_transport() {
assert!(runtime_instance_change_requires_replacement(true, false));
assert!(!runtime_instance_change_requires_replacement(false, false));
}
#[test]
fn ticket_change_clears_transient_retry_without_retaining_ticket_material() {
let mut last_known = std::collections::HashMap::new();
let mut last_attempt_at = std::collections::HashMap::new();
let mut failure_count = std::collections::HashMap::new();
let mut failure_backoff_until = std::collections::HashMap::new();
assert!(!reconcile_remote_ticket_retry_state(
&mut last_known,
"device-1",
"ticket-1.token-1",
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
));
last_attempt_at.insert("device-1".to_string(), 100);
failure_count.insert("device-1".to_string(), 4);
failure_backoff_until.insert("device-1".to_string(), 8_100);
assert!(!reconcile_remote_ticket_retry_state(
&mut last_known,
"device-1",
"ticket-1.token-1",
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
));
assert_eq!(last_attempt_at.get("device-1"), Some(&100));
assert_eq!(failure_count.get("device-1"), Some(&4));
assert_eq!(failure_backoff_until.get("device-1"), Some(&8_100));
assert!(reconcile_remote_ticket_retry_state(
&mut last_known,
"device-1",
"ticket-1.token-2",
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
));
assert!(!last_attempt_at.contains_key("device-1"));
assert!(!failure_count.contains_key("device-1"));
assert!(!failure_backoff_until.contains_key("device-1"));
let stored = last_known
.get("device-1::ticket")
.expect("ticket fingerprint");
assert_eq!(stored.len(), 16);
assert!(!stored.contains("token-2"));
}
#[test]
fn replacement_window_suppresses_auto_connect_for_active_replacement_churn() {
let snapshot = replacement_snapshot(Some("incoming-replacement-churn"));
assert!(should_suppress_auto_connect_for_replacement_window(
&snapshot,
snapshot.last_lifecycle_transition_at_ms + 1_000,
false,
));
assert!(
!should_suppress_auto_connect_for_replacement_window(
&snapshot,
snapshot.last_lifecycle_transition_at_ms + 1_000,
true,
),
"replacement suppression must not bypass positive liveness checks for an owned transport"
);
}
#[test]
fn transient_disconnect_replacement_window_allows_auto_connect_redial() {
let snapshot = replacement_snapshot(Some("transient-disconnect"));
assert!(!should_suppress_auto_connect_for_replacement_window(
&snapshot,
snapshot.last_lifecycle_transition_at_ms + 1_000,
false,
));
}
#[test]
fn expired_replacement_window_allows_auto_connect_redial() {
let snapshot = replacement_snapshot(Some("incoming-replacement-churn"));
assert!(!should_suppress_auto_connect_for_replacement_window(
&snapshot,
snapshot.last_lifecycle_transition_at_ms + MANAGED_SETTLE_DEADLINE_MS,
false,
));
}
#[test]
fn remote_exclusion_matches_canonical_device_id_case_insensitively() {
let excluded_peers = vec![" DEVICE-LOCAL ".to_string()];
assert!(remote_excludes_local_device(
&excluded_peers,
"device-local"
));
assert!(!remote_excludes_local_device(
&excluded_peers,
"device-other"
));
assert!(!remote_excludes_local_device(&excluded_peers, ""));
}
#[test]
fn browser_desired_peer_mailbox_fences_stale_revisions_and_replaces_state() {
let mut mailbox = BrowserDesiredPeerMailbox::default();
let peer = BrowserDesiredPeer {
device_id: " device-1 ".to_string(),
node_id: Some(" node-1 ".to_string()),
ticket: Some(" ticket-1 ".to_string()),
online: true,
session_id: Some(" runtime-1 ".to_string()),
excluded_peers: vec![" DEVICE-LOCAL ".to_string()],
};
assert!(mailbox.accept_with_withdrawn(2, vec![peer]).is_some());
assert!(mailbox.accept_with_withdrawn(1, Vec::new()).is_none());
assert!(mailbox.accept_with_withdrawn(2, Vec::new()).is_none());
assert_eq!(mailbox.revision, 2);
let stored = mailbox.peers.get("device-1").expect("desired peer");
assert_eq!(stored.node_id.as_deref(), Some("node-1"));
assert_eq!(stored.excluded_peers, vec!["device-local"]);
let withdrawn = mailbox
.accept_with_withdrawn(3, Vec::new())
.expect("newer revision is accepted");
assert_eq!(withdrawn.len(), 1);
assert_eq!(withdrawn[0].device_id, "device-1");
assert!(mailbox.peers.is_empty());
}
#[test]
fn browser_desired_peer_mailbox_withdraws_live_route_when_device_goes_offline() {
let mut mailbox = BrowserDesiredPeerMailbox::default();
let online = BrowserDesiredPeer {
device_id: "device-1".to_string(),
node_id: Some("node-1".to_string()),
ticket: Some("ticket-1".to_string()),
online: true,
session_id: Some("runtime-1".to_string()),
excluded_peers: Vec::new(),
};
assert!(mailbox
.accept_with_withdrawn(1, vec![online.clone()])
.expect("initial revision is accepted")
.is_empty());
let mut offline = online.clone();
offline.online = false;
let withdrawn = mailbox
.accept_with_withdrawn(2, vec![offline])
.expect("offline revision is accepted");
assert_eq!(withdrawn, vec![online]);
assert!(mailbox.peers.contains_key("device-1"));
assert!(!mailbox.peers["device-1"].online);
assert!(mailbox
.accept_with_withdrawn(3, vec![mailbox.peers["device-1"].clone()])
.expect("repeated offline revision is accepted")
.is_empty());
}
#[test]
fn browser_peer_evidence_fingerprints_ticket_instead_of_retaining_it() {
let peer = BrowserDesiredPeer {
device_id: "device-1".to_string(),
node_id: Some("node-1".to_string()),
ticket: Some("iroh-ticket.sensitive-token-material".to_string()),
online: true,
session_id: Some("runtime-1".to_string()),
excluded_peers: Vec::new(),
};
let evidence = browser_peer_evidence(&peer);
assert!(evidence.starts_with("node-1|"));
assert!(!evidence.contains("sensitive-token-material"));
}
#[test]
fn desired_peer_attempt_survives_unrelated_mailbox_revision() {
let mut mailbox = BrowserDesiredPeerMailbox::default();
let tracked = BrowserDesiredPeer {
device_id: "device-tracked".to_string(),
node_id: Some("node-tracked".to_string()),
ticket: Some("ticket-tracked".to_string()),
online: true,
session_id: Some("runtime-tracked".to_string()),
excluded_peers: Vec::new(),
};
let unrelated = BrowserDesiredPeer {
device_id: "device-unrelated".to_string(),
node_id: Some("node-unrelated".to_string()),
ticket: Some("ticket-unrelated".to_string()),
online: true,
session_id: Some("runtime-unrelated".to_string()),
excluded_peers: Vec::new(),
};
assert!(mailbox.accept(1, vec![tracked.clone()]));
assert!(desired_peer_evidence_is_current(&mailbox, &tracked));
assert!(mailbox.accept(2, vec![tracked.clone(), unrelated.clone()]));
assert!(desired_peer_evidence_is_current(&mailbox, &tracked));
let mut changed = tracked.clone();
changed.session_id = Some("runtime-replacement".to_string());
assert!(mailbox.accept(3, vec![changed, unrelated.clone()]));
assert!(!desired_peer_evidence_is_current(&mailbox, &tracked));
assert!(mailbox.accept(4, vec![unrelated]));
assert!(!desired_peer_evidence_is_current(&mailbox, &tracked));
}
#[test]
fn browser_desired_peer_tie_break_has_one_initial_dialer_and_bounded_takeover() {
assert_eq!(
browser_desired_peer_decision("node-z", "node-a", false, 1_000, None),
BrowserDesiredPeerDecision::Dial
);
assert_eq!(
browser_desired_peer_decision("node-a", "node-z", false, 1_000, None),
BrowserDesiredPeerDecision::WaitUntil(4_000)
);
assert_eq!(
browser_desired_peer_decision("node-a", "node-z", false, 4_000, Some(4_000)),
BrowserDesiredPeerDecision::Dial
);
assert_eq!(
browser_desired_peer_decision("node-a", "node-z", true, 1_000, None),
BrowserDesiredPeerDecision::ObserveExisting
);
}
#[test]
fn browser_retry_deadline_only_replaces_with_earlier_work_or_cancellation() {
assert!(!should_replace_browser_retry_deadline(None, None));
assert!(should_replace_browser_retry_deadline(None, Some(5_000)));
assert!(!should_replace_browser_retry_deadline(
Some(5_000),
Some(6_000)
));
assert!(should_replace_browser_retry_deadline(
Some(5_000),
Some(4_000)
));
assert!(should_replace_browser_retry_deadline(Some(5_000), None));
}
#[test]
fn browser_failure_backoff_is_bounded_for_recovery() {
assert_eq!(browser_auto_connect_failure_backoff_ms(1), 1_000);
assert_eq!(browser_auto_connect_failure_backoff_ms(2), 2_000);
assert_eq!(browser_auto_connect_failure_backoff_ms(10), 8_000);
}
#[test]
fn reconcile_deadlines_replace_expired_non_initiator_wait_with_failure_backoff() {
let (retry_deadline, non_initiator_deadline) =
browser_reconcile_deadlines(BrowserPeerReconcileOutcome::Retry, Some(9_000));
assert_eq!(retry_deadline, Some(9_000));
assert_eq!(non_initiator_deadline, None);
let (retry_deadline, non_initiator_deadline) =
browser_reconcile_deadlines(BrowserPeerReconcileOutcome::WaitUntil(4_000), None);
assert_eq!(retry_deadline, None);
assert_eq!(non_initiator_deadline, Some(4_000));
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn duplicate_external_auto_connect_preserves_gateway_mailbox_generation() {
let client = std::sync::Arc::new(Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
"test-app".to_string(),
Box::new(|| None),
));
client
.start_external_auto_connect("test-user".to_string(), "local-device".to_string())
.await
.expect("external actor starts");
assert!(client
.submit_external_desired_peers(7, "[]")
.await
.expect("gateway revision is accepted"));
let before = client.external_desired_peer_debug_snapshot().await;
client
.start_external_auto_connect("test-user".to_string(), "local-device".to_string())
.await
.expect("duplicate start is idempotent");
let after = client.external_desired_peer_debug_snapshot().await;
assert_eq!(after["generation"], before["generation"]);
assert_eq!(after["revision"], serde_json::json!(7));
assert_eq!(after["active"], serde_json::json!(true));
}
}
fn should_suppress_auto_connect_for_replacement_window(
snapshot: &ManagedConnectionHealthSnapshot,
now_ms: i64,
transport_connected: bool,
) -> bool {
if !snapshot.replacement_pending || transport_connected {
return false;
}
if now_ms.saturating_sub(snapshot.last_lifecycle_transition_at_ms) >= MANAGED_SETTLE_DEADLINE_MS
{
return false;
}
!matches!(
snapshot.last_reconnect_reason.as_deref(),
Some("transient-disconnect")
)
}
impl Client {
async fn retire_desired_peer_connections(
&self,
remote_device_id: &str,
remote_node_id: Option<&str>,
reason: &str,
) {
let mut records = self.resolve_peer_connection_records(remote_device_id).await;
if records.is_empty() {
if let Some(node_id) = remote_node_id
.map(str::trim)
.filter(|value| !value.is_empty())
{
records = self.resolve_peer_connection_records(node_id).await;
}
}
for record in records {
let endpoint_id = record
.endpoint_id
.as_deref()
.or(record.node_id.as_deref())
.and_then(|value| value.trim().parse::<iroh::EndpointId>().ok());
let disconnected = if let Some(endpoint_id) = endpoint_id {
self.disconnect_with_reason(endpoint_id, reason)
.await
.is_ok()
} else {
false
};
if !disconnected {
self.retire_managed_connection(&record.connection_id, Some(reason.to_string()))
.await;
}
println!(
"[pluto-rtc][auto-connect] retired undesired connection remote_device_id={} connection_id={} reason={} transport_closed={}",
remote_device_id, record.connection_id, reason, disconnected
);
}
}
async fn retire_withdrawn_desired_peer_connections(
&self,
remote_device_id: &str,
remote_node_id: Option<&str>,
) {
self.retire_desired_peer_connections(
remote_device_id,
remote_node_id,
crate::lifecycle_reason::REASON_DESIRED_PEER_WITHDRAWN,
)
.await;
}
pub async fn is_connection_transport_alive(&self, endpoint_id: iroh::EndpointId) -> bool {
if !self.is_connected(endpoint_id).await {
return false;
}
#[cfg(not(target_arch = "wasm32"))]
{
let Some(connection) = self.get_connection(endpoint_id).await else {
return false;
};
return connection.close_reason().is_none();
}
#[cfg(target_arch = "wasm32")]
{
let Some(connection) = self.get_connection(endpoint_id).await else {
return false;
};
connection.close_reason().is_none()
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn is_managed_connection_transport_alive(&self, connection_id: &str) -> bool {
let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return false;
};
let Some(node_id) = record
.node_id
.as_deref()
.map(str::trim)
.filter(|v| !v.is_empty())
else {
return false;
};
let Ok(endpoint_id) = node_id.parse::<iroh::EndpointId>() else {
return false;
};
self.is_connection_transport_alive(endpoint_id).await
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn validate_active_connection(
&self,
endpoint_id: iroh::EndpointId,
) -> Option<crate::heartbeat::IrohConnectionProbe> {
if !self.is_connection_transport_alive(endpoint_id).await {
return None;
}
let node = self.iroh_node.read().await.as_ref().cloned()?;
node.probe_connection(endpoint_id, std::time::Duration::from_secs(4))
.await
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn should_force_network_change_after_connect_failures(
&self,
remote_device_id: &str,
) -> bool {
for record in self.connection_manager.list_active().await {
if !matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) {
continue;
}
let Some(record_device_id) = record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
else {
continue;
};
if record_device_id == remote_device_id {
continue;
}
println!(
"[pluto-rtc][auto-connect] suppress network-change recovery remote_device_id={} because healthy peer remains connected peer_device_id={} connection_id={}",
remote_device_id,
record_device_id,
record.connection_id
);
return false;
}
true
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn try_auto_connect_device(
self: &Arc<Self>,
user_id: &str,
device: crate::signaling::Device,
local_device_id: &str,
last_attempt_at: &mut std::collections::HashMap<String, i64>,
failure_count: &mut std::collections::HashMap<String, u8>,
failure_backoff_until: &mut std::collections::HashMap<String, i64>,
retry_throttle_ms: i64,
non_initiator_wait_started_at: &mut std::collections::HashMap<String, i64>,
non_initiator_escalation_count: &mut std::collections::HashMap<String, u8>,
initiator_grace_ms: i64,
skip_log_at: &mut std::collections::HashMap<String, i64>,
connected_health_checked_at: &mut std::collections::HashMap<String, i64>,
connected_health_state: &mut std::collections::HashMap<String, bool>,
connected_health_failures: &mut std::collections::HashMap<String, u8>,
connected_health_probe_ms: i64,
connected_health_failure_threshold: u8,
last_presence_republish_at_ms: &mut i64,
presence_republish_interval_ms: i64,
last_network_change_at_ms: &mut i64,
network_change_recovery_interval_ms: i64,
network_change_failure_threshold: u8,
last_known_node_id: &mut std::collections::HashMap<String, String>,
) {
if self.is_app_backgrounded() {
return;
}
#[cfg(target_arch = "wasm32")]
let _ = (
user_id,
last_presence_republish_at_ms,
presence_republish_interval_ms,
last_network_change_at_ms,
network_change_recovery_interval_ms,
network_change_failure_threshold,
);
let remote_device_id = device.device_id.clone();
let now = now_millis_i64();
let skip_log_interval_ms = 2_000;
let mut log_skip = |reason: &str| {
let is_noise_reason = matches!(
reason,
"already-connected"
| "tie-break-not-initiator-waiting"
| "retry-throttle"
| "self-ticket"
| "auto-connect-excluded"
| "remote-excluded-local-device"
);
if is_noise_reason && !auto_connect_verbose() {
return;
}
let skip_key = format!("{}::{}", remote_device_id, reason);
let should_log = match skip_log_at.get(&skip_key) {
Some(last_at) => now.saturating_sub(*last_at) >= skip_log_interval_ms,
None => true,
};
if !should_log {
return;
}
skip_log_at.insert(skip_key, now);
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] skip remote_device_id={} local_device_id={} reason={}",
remote_device_id, local_device_id, reason
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] skip remote_device_id={} local_device_id={} reason={}",
remote_device_id, local_device_id, reason
)));
};
if !device.online {
last_attempt_at.remove(&remote_device_id);
non_initiator_wait_started_at.remove(&remote_device_id);
non_initiator_escalation_count.remove(&remote_device_id);
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
return;
}
if remote_device_id == local_device_id {
last_attempt_at.remove(&remote_device_id);
non_initiator_wait_started_at.remove(&remote_device_id);
non_initiator_escalation_count.remove(&remote_device_id);
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
return;
}
if remote_excludes_local_device(&device.excluded_peers, local_device_id) {
log_skip("remote-excluded-local-device");
return;
}
if self.is_auto_connect_peer_excluded(&remote_device_id, device.node_id.as_deref()) {
log_skip("auto-connect-excluded");
return;
}
let ticket_str = match device.ticket.as_deref() {
Some(ticket) if !ticket.trim().is_empty() => ticket,
_ => {
log_skip("missing-ticket");
return;
}
};
let (iroh_ticket, token_suffix) =
crate::session_token::split_compound_ticket(ticket_str.trim());
let extracted_token_payload = token_suffix.and_then(|suffix| {
crate::session_token::decode_token_payload_for_ticket(iroh_ticket, suffix)
});
let extracted_token = extracted_token_payload
.as_ref()
.map(|payload| payload.token.clone());
let extracted_token_suffix = token_suffix
.filter(|_| extracted_token_payload.is_some())
.map(ToOwned::to_owned);
let endpoint_addr = match parse_endpoint_ticket(iroh_ticket) {
Ok(addr) => addr,
Err(error) => {
log_skip("invalid-ticket");
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[pluto-rtc][auto-connect] invalid ticket remote_device_id={} error={}",
remote_device_id, error
);
#[cfg(target_arch = "wasm32")]
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] invalid ticket remote_device_id={} error={}",
remote_device_id, error
)));
return;
}
};
let endpoint_id = endpoint_addr.id;
let node_id_str = endpoint_id.to_string();
#[cfg(not(target_arch = "wasm32"))]
self.update_native_route_repair_credential(
&node_id_str,
Some(&remote_device_id),
extracted_token.as_deref(),
extracted_token_suffix.as_deref(),
);
let runtime_instance_changed = observe_remote_runtime_instance(
last_known_node_id,
&remote_device_id,
device.session_id.as_deref(),
);
let ticket_changed = reconcile_remote_ticket_retry_state(
last_known_node_id,
&remote_device_id,
ticket_str,
last_attempt_at,
failure_count,
failure_backoff_until,
);
if ticket_changed {
println!(
"[pluto-rtc][auto-connect] peer ticket changed remote_device_id={} node_id={} - clearing stale retry deadline",
remote_device_id, node_id_str
);
}
{
let identity_changed = last_known_node_id
.get(&remote_device_id)
.map(|prev| prev.as_str() != node_id_str)
.unwrap_or(false);
last_known_node_id.insert(remote_device_id.clone(), node_id_str.clone());
if identity_changed {
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] peer identity changed remote_device_id={} new_node_id={} — resetting state for fast reconnect",
remote_device_id, node_id_str
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] peer identity changed remote_device_id={} new_node_id={} — resetting state for fast reconnect",
remote_device_id, node_id_str
)));
last_attempt_at.remove(&remote_device_id);
non_initiator_wait_started_at.remove(&remote_device_id);
non_initiator_escalation_count.remove(&remote_device_id);
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
}
}
let local_node_id = match self.current_node_id().await {
Some(id) => id,
None => {
log_skip("local-node-not-ready");
return;
}
};
if local_node_id == node_id_str {
log_skip("self-ticket");
return;
}
if runtime_instance_changed {
let current_transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
#[cfg(openrtc_iroh_retire_remote_connections_api)]
if let Some(transport_stable_id) = current_transport_stable_id {
let endpoint = { self.iroh_endpoint.read().await.clone() };
if let Some(endpoint) = endpoint {
let retired = endpoint
.retire_remote_connections_except(
endpoint_id,
transport_stable_id as usize,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT
.as_bytes(),
)
.await;
if retired > 0 {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - retired {} stale shared-endpoint protocol connection(s), preserved transport_stable_id={}",
remote_device_id, node_id_str, retired, transport_stable_id
);
}
}
}
let active_transport_responsive = if current_transport_stable_id.is_some() {
self.validate_active_connection(endpoint_id)
.await
.is_some_and(|probe| probe.responsive)
} else {
false
};
if runtime_instance_change_requires_replacement(
runtime_instance_changed,
active_transport_responsive,
) {
if let Some(transport_stable_id) = current_transport_stable_id {
match self
.retire_transport_generation_after_failed_probe(
&Self::deterministic_connection_id(&local_node_id, &node_id_str),
endpoint_id,
transport_stable_id,
connected_health_probe_ms,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await
{
FailedProbeRetirementOutcome::Retired => {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - replaced unresponsive physical transport",
remote_device_id, node_id_str
);
}
FailedProbeRetirementOutcome::DeferredAdmissionInFlight => {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - deferred replacement while admission is in flight",
remote_device_id, node_id_str
);
}
FailedProbeRetirementOutcome::CancelledCurrentGenerationActivity => {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - preserved generation with recent protocol activity",
remote_device_id, node_id_str
);
}
FailedProbeRetirementOutcome::GenerationChanged => {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - ignored stale generation",
remote_device_id, node_id_str
);
}
}
}
} else {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - preserving responsive physical transport",
remote_device_id, node_id_str
);
}
last_attempt_at.remove(&remote_device_id);
non_initiator_wait_started_at.remove(&remote_device_id);
non_initiator_escalation_count.remove(&remote_device_id);
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
}
if let Some(discovered_node_id) = device.node_id.as_deref() {
if discovered_node_id != node_id_str {
log_skip("ticket-node-mismatch-discovery");
eprintln!(
"[pluto-rtc][auto-connect] rejected inconsistent discovery identity remote_device_id={} discovered_node_id={} ticket_node_id={}",
remote_device_id, discovered_node_id, node_id_str
);
return;
}
}
self.reconcile_authoritative_device_node(&remote_device_id, &node_id_str)
.await;
#[cfg(feature = "transport-lan")]
if self.is_peer_locally_reachable(&node_id_str).await {
println!(
"[pluto-rtc][auto-connect] LAN peer reachable remote_device_id={} node_id={}",
remote_device_id, node_id_str
);
}
let connection_id = Self::deterministic_connection_id(&local_node_id, &node_id_str);
let managed_health_snapshot = self.managed_connection_health(&connection_id).await;
#[cfg(not(target_arch = "wasm32"))]
{
if let Some(existing) = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
let existing_device_id = existing.device_id.as_deref().map(str::trim);
if existing_device_id != Some(remote_device_id.as_str()) {
self.connection_manager
.set_device_id(&connection_id, remote_device_id.clone())
.await;
println!(
"[pluto-rtc][auto-connect] backfilled device_id hint connection_id={} node_id={} device_id={}",
connection_id, node_id_str, remote_device_id
);
}
}
}
let transport_connected = self.is_connected(endpoint_id).await;
if let Some(snapshot) = managed_health_snapshot.as_ref() {
if should_suppress_auto_connect_for_replacement_window(
&snapshot,
now,
transport_connected,
) {
log_skip("replacement-pending");
return;
}
}
#[cfg(target_arch = "wasm32")]
if !transport_connected {
if let Some(existing) = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
if matches!(
existing.state,
crate::connection_manager::ConnectionState::Connected
) {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] transport mismatch remote_device_id={} remote_node_id={} connection_id={} record_stable_id={:?} record_generation={} state={:?} status_reason={:?}",
remote_device_id,
node_id_str,
connection_id,
existing.transport_stable_id,
existing.transport_generation,
existing.state,
existing.status_reason
)));
}
}
}
let esc_count = non_initiator_escalation_count
.get(&remote_device_id)
.copied()
.unwrap_or(0);
let effective_grace_ms = if local_node_id <= node_id_str {
non_initiator_escalation_grace_ms(esc_count)
} else {
initiator_grace_ms
};
let waited_ms = if local_node_id <= node_id_str {
let wait_started_at = non_initiator_wait_started_at
.entry(remote_device_id.clone())
.or_insert(now);
now.saturating_sub(*wait_started_at)
} else {
0
};
match auto_connect_tie_break_decision(
&local_node_id,
&node_id_str,
transport_connected,
waited_ms,
effective_grace_ms,
) {
AutoConnectTieBreakDecision::ObserveConnectedTransport => {
non_initiator_wait_started_at.remove(&remote_device_id);
non_initiator_escalation_count.remove(&remote_device_id);
}
AutoConnectTieBreakDecision::WaitForInitiator => {
if waited_ms < 1_000 {
log_skip("tie-break-not-initiator-waiting");
}
return;
}
AutoConnectTieBreakDecision::ActAsInitiator => {
if local_node_id <= node_id_str {
let new_esc_count = esc_count.saturating_add(1);
non_initiator_escalation_count.insert(remote_device_id.clone(), new_esc_count);
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] non-initiator escalated after {}ms grace remote_device_id={} local_device_id={} escalation_count={}",
waited_ms, remote_device_id, local_device_id, new_esc_count
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] non-initiator escalated after {}ms grace remote_device_id={} local_device_id={} escalation_count={}",
waited_ms, remote_device_id, local_device_id, new_esc_count
)));
}
non_initiator_wait_started_at.remove(&remote_device_id);
}
}
self.reconcile_authoritative_device_node(&remote_device_id, &node_id_str)
.await;
if transport_connected {
connected_health_checked_at
.entry(remote_device_id.clone())
.or_insert(now);
connected_health_state
.entry(remote_device_id.clone())
.or_insert(true);
let should_probe = match connected_health_checked_at.get(&remote_device_id) {
Some(last_probe_at) => {
now.saturating_sub(*last_probe_at) >= connected_health_probe_ms
}
None => false,
};
let mut healthy = connected_health_state
.get(&remote_device_id)
.copied()
.unwrap_or(true);
let mut probed_transport_stable_id = None;
if should_probe {
let active_probe = self.validate_active_connection(endpoint_id).await;
probed_transport_stable_id =
active_probe.as_ref().map(|probe| probe.transport_stable_id);
let active_probe_healthy =
active_probe.as_ref().is_some_and(|probe| probe.responsive);
let passive_transport_alive = self.is_connection_transport_alive(endpoint_id).await;
let independent_route_ready = self
.native_active_independent_route_ready(&connection_id)
.await;
let recent_generation_activity = active_probe.as_ref().is_some_and(|probe| {
generation_activity_is_recent(
probe.last_inbound_activity_age,
connected_health_probe_ms,
) || generation_activity_is_recent(
self.native_transport_protocol_activity_age(
&connection_id,
probe.transport_stable_id,
),
connected_health_probe_ms,
)
});
let probe_decision = active_connection_probe_decision_with_independent_route(
active_probe_healthy,
passive_transport_alive,
independent_route_ready,
recent_generation_activity,
);
let current_failures = connected_health_failures
.get(&remote_device_id)
.copied()
.unwrap_or(0);
let failures =
next_active_connection_probe_failure_count(probe_decision, current_failures);
healthy = active_connection_probe_is_healthy(
probe_decision,
failures,
connected_health_failure_threshold,
);
connected_health_checked_at.insert(remote_device_id.clone(), now);
connected_health_state.insert(remote_device_id.clone(), healthy);
if matches!(probe_decision, ActiveConnectionProbeDecision::Healthy) {
connected_health_failures.remove(&remote_device_id);
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
} else {
connected_health_failures.insert(remote_device_id.clone(), failures);
if matches!(
probe_decision,
ActiveConnectionProbeDecision::PreserveLiveTransport
) && healthy
{
log_skip("active-probe-missed-preserving-live-transport");
}
}
}
if !healthy {
if !should_probe {
log_skip("awaiting-current-generation-health-probe");
return;
}
if self.is_app_backgrounded() {
log_skip("backgrounded-after-health-probe");
return;
}
#[cfg(not(target_arch = "wasm32"))]
{
let failures = connected_health_failures
.get(&remote_device_id)
.copied()
.unwrap_or(1);
if failures < connected_health_failure_threshold {
log_skip("active-connection-health-probe-failed");
return;
}
let Some(transport_stable_id) = probed_transport_stable_id else {
log_skip("active-probe-generation-disappeared");
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
return;
};
match self
.retire_transport_generation_after_failed_probe(
&connection_id,
endpoint_id,
transport_stable_id,
connected_health_probe_ms,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await
{
FailedProbeRetirementOutcome::Retired => {
log_skip("stale-active-connection-reconnect");
}
FailedProbeRetirementOutcome::DeferredAdmissionInFlight => {
log_skip("active-probe-retirement-deferred-admission-in-flight");
connected_health_checked_at.insert(remote_device_id.clone(), now);
connected_health_state.insert(remote_device_id.clone(), true);
connected_health_failures.remove(&remote_device_id);
return;
}
FailedProbeRetirementOutcome::CancelledCurrentGenerationActivity => {
log_skip(
"active-probe-retirement-cancelled-by-current-generation-activity",
);
connected_health_checked_at.insert(remote_device_id.clone(), now);
connected_health_state.insert(remote_device_id.clone(), true);
connected_health_failures.remove(&remote_device_id);
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
return;
}
FailedProbeRetirementOutcome::GenerationChanged => {
log_skip("active-probe-generation-disappeared");
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
return;
}
}
self.maybe_refresh_gateway_presence_for_auto_connect(
user_id,
local_device_id,
"stale-active-connection-reconnect",
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
self.reconcile_authoritative_device_node(&remote_device_id, &node_id_str)
.await;
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
if failures >= network_change_failure_threshold {
self.maybe_force_network_change_for_auto_connect(
user_id,
local_device_id,
"stale-active-connection",
last_network_change_at_ms,
network_change_recovery_interval_ms,
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
}
}
#[cfg(target_arch = "wasm32")]
{
let failures = connected_health_failures
.get(&remote_device_id)
.copied()
.unwrap_or(1);
if failures < connected_health_failure_threshold {
println!(
"[pluto-rtc][auto-connect][health] remote_device_id={} local_device_id={} remote_node_id={} healthy=false failures={} threshold={} action=wait",
remote_device_id,
local_device_id,
node_id_str,
failures,
connected_health_failure_threshold
);
log_skip("active-connection-health-probe-failed");
return;
}
println!(
"[pluto-rtc][auto-connect][health] remote_device_id={} local_device_id={} remote_node_id={} healthy=false failures={} threshold={} action=disconnect-and-reconnect",
remote_device_id,
local_device_id,
node_id_str,
failures,
connected_health_failure_threshold
);
log_skip("stale-active-connection-reconnect");
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await;
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
}
} else {
if self.is_auto_connect_peer_excluded(&remote_device_id, Some(&node_id_str)) {
log_skip("auto-connect-excluded-after-health-probe");
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_MANUAL_DISCONNECT,
)
.await;
return;
}
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
let current_transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(node_id_str.clone()),
current_transport_stable_id,
Some("health-check".to_string()),
)
.await;
#[cfg(target_arch = "wasm32")]
let _ = self
.report_transport_status_for_current_generation(
&connection_id,
"iroh-relay",
None,
)
.await;
#[cfg(not(target_arch = "wasm32"))]
let has_remote_admission_proof =
current_transport_stable_id.is_some_and(|transport_stable_id| {
self.remote_session_token_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
#[cfg(target_arch = "wasm32")]
let has_remote_admission_proof = if let Some(token) = extracted_token.as_deref() {
self.has_current_remote_session_admission_proof(
&connection_id,
endpoint_id,
token,
)
.await
} else {
false
};
#[cfg(not(target_arch = "wasm32"))]
let has_current_inbound_admission_proof =
current_transport_stable_id.is_some_and(|transport_stable_id| {
self.inbound_session_token_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
#[cfg(not(target_arch = "wasm32"))]
let peer_requested_reciprocal =
self.has_pending_reciprocal_session_admission_request(&connection_id);
#[cfg(not(target_arch = "wasm32"))]
if has_remote_admission_proof
&& has_current_inbound_admission_proof
&& peer_requested_reciprocal
{
self.clear_reciprocal_session_admission_request(&connection_id);
}
#[cfg(not(target_arch = "wasm32"))]
let presentation_decision = session_admission_presentation_decision(
extracted_token.is_some(),
has_remote_admission_proof,
has_current_inbound_admission_proof,
peer_requested_reciprocal,
local_node_id > node_id_str,
);
#[cfg(target_arch = "wasm32")]
let presentation_decision = session_admission_presentation_decision(
extracted_token.is_some(),
has_remote_admission_proof,
true,
false,
true,
);
if presentation_decision.should_present {
if let Some(backoff_until) =
failure_backoff_until.get(&remote_device_id).copied()
{
if now < backoff_until {
log_skip("session-admission-backoff");
return;
}
failure_backoff_until.remove(&remote_device_id);
}
let token = extracted_token
.as_ref()
.expect("representation requires an extracted token");
println!(
"[pluto-rtc][auto-connect][session-admission] re-presenting token on live transport connection_id={} remote_device_id={} remote_node_id={} token_fp={}",
connection_id,
remote_device_id,
node_id_str,
super::core_impl::log_fingerprint(token.as_str())
);
#[cfg(not(target_arch = "wasm32"))]
let presentation_result = self
.present_and_accept_session_token_for_route_repair(
endpoint_id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
Some(remote_device_id.clone()),
Some(local_device_id.to_string()),
has_remote_admission_proof,
presentation_decision.request_reciprocal,
)
.await;
#[cfg(target_arch = "wasm32")]
let presentation_result = self
.present_and_accept_session_token_with_local_claim(
endpoint_id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
Some(remote_device_id.clone()),
Some(local_device_id.to_string()),
)
.await;
match presentation_result {
Ok(approval_scope) => {
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
println!(
"[pluto-rtc][auto-connect][session-admission] live transport admitted connection_id={} remote_device_id={} remote_node_id={} scope={} token_fp={}",
connection_id,
remote_device_id,
node_id_str,
approval_scope,
super::core_impl::log_fingerprint(token.as_str())
);
}
Err(error) => {
let active_transport_stable_id =
self.get_connection(endpoint_id).await.map(|connection| {
crate::transport_generation::for_connection(&connection)
});
if current_transport_stable_id != active_transport_stable_id {
println!(
"[pluto-rtc][auto-connect][session-admission] stale presentation failure ignored because replacement already owns endpoint connection_id={} remote_device_id={} attempted_transport_stable_id={:?} active_transport_stable_id={:?}",
connection_id,
remote_device_id,
current_transport_stable_id,
active_transport_stable_id
);
return;
}
#[cfg(not(target_arch = "wasm32"))]
let current_generation_is_routable = active_transport_stable_id
.is_some_and(|transport_stable_id| {
self.native_application_stream_admitted_for_transport(
&connection_id,
transport_stable_id,
)
});
#[cfg(target_arch = "wasm32")]
let current_generation_is_routable = false;
if current_generation_is_routable {
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
connected_health_failures.remove(&remote_device_id);
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
println!(
"[pluto-rtc][auto-connect][session-admission] ignored losing presentation failure because another exchange made the generation routable connection_id={} remote_device_id={} transport_stable_id={:?} error={}",
connection_id,
remote_device_id,
active_transport_stable_id,
error,
);
return;
}
if session_token_presentation_was_duplicate(&error) {
println!(
"[pluto-rtc][auto-connect][session-admission] ignored duplicate presentation owned by canonical exchange connection_id={} remote_device_id={} transport_stable_id={:?}",
connection_id,
remote_device_id,
active_transport_stable_id,
);
return;
}
if self
.session_token_registry
.is_session_token_presentation_in_flight(&connection_id)
{
println!(
"[pluto-rtc][auto-connect][session-admission] ignored duplicate presentation while canonical exchange remains in flight connection_id={} remote_device_id={} transport_stable_id={:?}",
connection_id,
remote_device_id,
active_transport_stable_id,
);
return;
}
if session_token_presentation_failure_action(&error)
== SessionTokenPresentationFailureAction::RejectAndDisconnect
{
let rejection_reason =
format!("session-token-presentation-failed: {}", error);
#[cfg(all(
not(target_arch = "wasm32"),
feature = "transport-webrtc"
))]
self.suppress_native_webrtc_restarts(
&connection_id,
Self::NATIVE_WEBRTC_ADMISSION_FAILURE_BACKOFF_MS,
rejection_reason.clone(),
)
.await;
self.reject_session_connection(
&connection_id,
rejection_reason.as_str(),
);
self.connection_manager
.set_failed(&connection_id, Some(rejection_reason.clone()))
.await;
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_SESSION_ADMISSION_REJECTED,
)
.await;
self.forget_session_connection(&connection_id);
register_auto_connect_admission_rejection(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
eprintln!(
"[pluto-rtc][auto-connect][session-admission] live transport admission rejected connection_id={} remote_device_id={} remote_node_id={} token_fp={} error={}",
connection_id,
remote_device_id,
node_id_str,
super::core_impl::log_fingerprint(token.as_str()),
error
);
return;
}
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
let failures =
failure_count.get(&remote_device_id).copied().unwrap_or(1);
match retryable_session_admission_failure_decision(
failures,
connected_health_failure_threshold,
current_transport_stable_id,
active_transport_stable_id,
runtime_instance_changed,
) {
RetryableSessionAdmissionFailureDecision::PreserveCurrentTransport => {
eprintln!(
"[pluto-rtc][auto-connect][session-admission] transient presentation failure; preserving current generation for one bounded retry connection_id={} remote_device_id={} remote_node_id={} transport_stable_id={:?} failures={} threshold={} token_fp={} error={}",
connection_id,
remote_device_id,
node_id_str,
current_transport_stable_id,
failures,
connected_health_failure_threshold,
super::core_impl::log_fingerprint(token.as_str()),
error
);
}
RetryableSessionAdmissionFailureDecision::RetireCurrentTransport => {
eprintln!(
"[pluto-rtc][auto-connect][session-admission] repeated transient presentation failure; retiring unresponsive generation for canonical reconnect connection_id={} remote_device_id={} remote_node_id={} transport_stable_id={:?} failures={} threshold={} token_fp={} error={}",
connection_id,
remote_device_id,
node_id_str,
current_transport_stable_id,
failures,
connected_health_failure_threshold,
super::core_impl::log_fingerprint(token.as_str()),
error
);
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await;
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
}
RetryableSessionAdmissionFailureDecision::ReplacementAlreadyWon => {
println!(
"[pluto-rtc][auto-connect][session-admission] transient failure belongs to retired generation; replacement preserved connection_id={} remote_device_id={} attempted_transport_stable_id={:?} active_transport_stable_id={:?}",
connection_id,
remote_device_id,
current_transport_stable_id,
active_transport_stable_id
);
}
}
return;
}
}
}
#[cfg(not(target_arch = "wasm32"))]
{
if self
.promote_known_native_user_device_connection(&connection_id, &node_id_str)
.await
.is_none()
{
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
}
}
#[cfg(target_arch = "wasm32")]
{
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
}
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
non_initiator_escalation_count.remove(&remote_device_id);
log_skip("already-connected");
return;
}
} else {
connected_health_checked_at.remove(&remote_device_id);
connected_health_state.remove(&remote_device_id);
connected_health_failures.remove(&remote_device_id);
}
if let Some(backoff_until) = failure_backoff_until.get(&remote_device_id).copied() {
if now < backoff_until {
log_skip("failure-backoff");
return;
}
failure_backoff_until.remove(&remote_device_id);
}
let last = last_attempt_at.get(&remote_device_id).cloned().unwrap_or(0);
if now - last < retry_throttle_ms {
log_skip("retry-throttle");
return;
}
last_attempt_at.insert(remote_device_id.clone(), now);
let connect_mode = "ticket-endpoint-addr";
#[cfg(not(target_arch = "wasm32"))]
{
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
self.connection_manager.set_connecting(&connection_id).await;
println!(
"[pluto-rtc][auto-connect] attempt remote_device_id={} local_device_id={} local_node_id={} remote_node_id={} mode={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str,
connect_mode
);
println!("[pluto-rtc] auto-connecting to {}", remote_device_id);
let result = self.ensure_connected_addr(endpoint_id, endpoint_addr).await;
match result {
Ok(()) => {
if self.is_auto_connect_peer_excluded(&remote_device_id, Some(&node_id_str)) {
log_skip("auto-connect-excluded-after-dial");
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_MANUAL_DISCONNECT,
)
.await;
self.connection_manager
.set_closed(
&connection_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
return;
}
let active_probe = self.validate_active_connection(endpoint_id).await;
let active_probe_stable =
active_probe.as_ref().is_some_and(|probe| probe.responsive);
let transport_alive = self.is_connection_transport_alive(endpoint_id).await;
let probe_decision =
active_connection_probe_decision(active_probe_stable, transport_alive);
let stable = matches!(
probe_decision,
ActiveConnectionProbeDecision::Healthy
| ActiveConnectionProbeDecision::PreserveLiveTransport
);
if !stable {
if self.is_app_backgrounded() {
log_skip("backgrounded-after-stability-probe");
return;
}
self.connection_manager
.set_failed(
&connection_id,
Some("connection failed stability verification".to_string()),
)
.await;
if let Some(transport_stable_id) =
active_probe.as_ref().map(|probe| probe.transport_stable_id)
{
let _ = self
.disconnect_transport_generation_with_reason(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_CONNECTION_FAILED_STABILITY,
)
.await;
}
self.reconcile_authoritative_device_node(&remote_device_id, &node_id_str)
.await;
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
let current_failures =
failure_count.get(&remote_device_id).copied().unwrap_or(0);
if current_failures >= network_change_failure_threshold
&& self
.should_force_network_change_after_connect_failures(
&remote_device_id,
)
.await
{
self.maybe_force_network_change_for_auto_connect(
user_id,
local_device_id,
"unstable-connection",
last_network_change_at_ms,
network_change_recovery_interval_ms,
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
}
eprintln!(
"[pluto-rtc][auto-connect] unstable connection suppressed remote_device_id={} local_device_id={} local_node_id={} remote_node_id={} connection_id={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str,
connection_id
);
return;
}
if matches!(
probe_decision,
ActiveConnectionProbeDecision::PreserveLiveTransport
) {
log_skip("stability-probe-missed-preserving-live-transport");
}
let current_transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(node_id_str.clone()),
current_transport_stable_id,
Some(connect_mode.to_string()),
)
.await;
if let Some(ref token) = extracted_token {
let bilateral_user_device = extracted_token_payload
.as_ref()
.is_some_and(|payload| payload.scope.as_str() == "user-device");
let has_remote_admission_proof =
current_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 = current_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);
let presentation_decision = if bilateral_user_device {
session_admission_presentation_decision(
true,
has_remote_admission_proof,
has_current_inbound_admission_proof,
peer_requested_reciprocal,
local_node_id > node_id_str,
)
} else {
SessionAdmissionPresentationDecision {
should_present: !has_remote_admission_proof,
request_reciprocal: false,
}
};
if !presentation_decision.should_present {
println!(
"[pluto-rtc][auto-connect][session-admission] awaiting elected presenter on fresh transport connection_id={} remote_device_id={} remote_node_id={} local_node_id={}",
connection_id,
remote_device_id,
node_id_str,
local_node_id,
);
return;
}
match self
.present_and_accept_session_token_for_route_repair(
endpoint_id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
Some(remote_device_id.clone()),
Some(local_device_id.to_string()),
has_remote_admission_proof,
presentation_decision.request_reciprocal,
)
.await
{
Ok(approval_scope) => {
#[cfg(all(
not(target_arch = "wasm32"),
feature = "transport-webrtc"
))]
let _ = self
.clear_native_webrtc_suppression(&connection_id, None)
.await;
println!(
"[pluto-rtc][auto-connect] session-token presented and approved connection_id={} scope={} device_id={} token_fp={} local_admission=accepted",
connection_id,
approval_scope,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str())
);
}
Err(error) => {
if session_token_presentation_failure_action(&error)
== SessionTokenPresentationFailureAction::RetryOnNextTransport
{
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
eprintln!(
"[pluto-rtc][auto-connect] transient session-token presentation failure; preserving peer session for canonical transport retry connection_id={} device_id={} token_fp={} error={}",
connection_id,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str()),
error
);
return;
}
let rejection_reason =
format!("session-token-presentation-failed: {}", error);
#[cfg(all(
not(target_arch = "wasm32"),
feature = "transport-webrtc"
))]
self.suppress_native_webrtc_restarts(
&connection_id,
Self::NATIVE_WEBRTC_ADMISSION_FAILURE_BACKOFF_MS,
rejection_reason.clone(),
)
.await;
self.reject_session_connection(
&connection_id,
rejection_reason.as_str(),
);
self.connection_manager
.set_failed(&connection_id, Some(rejection_reason.clone()))
.await;
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_SESSION_ADMISSION_REJECTED,
)
.await;
self.reconcile_authoritative_device_node(
&remote_device_id,
&node_id_str,
)
.await;
self.forget_session_connection(&connection_id);
register_auto_connect_admission_rejection(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
eprintln!(
"[pluto-rtc][auto-connect] session-token presentation failed connection_id={} device_id={} token_fp={} error={}",
connection_id,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str()),
error
);
return;
}
}
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
if let Err(error) = self
.request_native_webrtc_recovery(
&connection_id,
Some(&node_id_str),
crate::native_webrtc_policy::NativeWebRTCRecoveryTrigger::Native(
crate::native_webrtc_policy::NativeWebRTCNativeTrigger::AutoConnectConnected,
),
crate::native_webrtc_policy::NativeWebRTCRecoveryOptions {
force_restart: true,
preferred_negotiation_id: None,
role_override: None,
},
)
.await
{
eprintln!(
"[NativeWebRTC] fallback trigger failed source=auto-connect-connected connection_id={} remote_node_id={} error={}",
connection_id,
node_id_str,
error
);
}
}
if let Some(outgoing_conn) = self.get_connection(endpoint_id).await {
self.spawn_iroh_path_watcher(
&connection_id,
&node_id_str,
&outgoing_conn,
false,
);
}
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
non_initiator_escalation_count.remove(&remote_device_id);
println!(
"[pluto-rtc][auto-connect] connected remote_device_id={} local_device_id={} local_node_id={} remote_node_id={} connection_id={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str,
connection_id
);
}
Err(e) => {
self.connection_manager
.set_failed(&connection_id, Some(e.to_string()))
.await;
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
let current_failures =
failure_count.get(&remote_device_id).copied().unwrap_or(0);
if current_failures >= network_change_failure_threshold
&& self
.should_force_network_change_after_connect_failures(&remote_device_id)
.await
{
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await;
self.reconcile_authoritative_device_node(&remote_device_id, &node_id_str)
.await;
self.maybe_force_network_change_for_auto_connect(
user_id,
local_device_id,
"repeated-connect-failures",
last_network_change_at_ms,
network_change_recovery_interval_ms,
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
}
self.maybe_refresh_gateway_presence_for_auto_connect(
user_id,
local_device_id,
"connect-failed",
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
eprintln!(
"[pluto-rtc] auto-connect failed to {}: {}",
remote_device_id, e
);
}
}
}
#[cfg(target_arch = "wasm32")]
{
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
self.connection_manager
.set_device_id(&connection_id, remote_device_id.clone())
.await;
self.connection_manager.set_connecting(&connection_id).await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] attempt remote_device_id={} local_device_id={} local_node_id={} remote_node_id={} mode={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str,
connect_mode
)));
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc] auto-connecting to {}",
remote_device_id
)));
let result = self.ensure_connected_addr(endpoint_id, endpoint_addr).await;
match result {
Ok(()) => {
let active_probe_stable = self.validate_active_connection(endpoint_id).await;
let transport_alive = self.is_connection_transport_alive(endpoint_id).await;
let probe_decision =
active_connection_probe_decision(active_probe_stable, transport_alive);
let stable = matches!(
probe_decision,
ActiveConnectionProbeDecision::Healthy
| ActiveConnectionProbeDecision::PreserveLiveTransport
);
if !stable {
if self.is_app_backgrounded() {
log_skip("backgrounded-after-stability-probe");
return;
}
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
self.connection_manager
.set_failed(
&connection_id,
Some("connection failed stability verification".to_string()),
)
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_AUTO_CONNECT_FAILURE,
)
.await;
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] unstable connection suppressed remote_device_id={} local_device_id={} local_node_id={} remote_node_id={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str
)));
return;
}
if matches!(
probe_decision,
ActiveConnectionProbeDecision::PreserveLiveTransport
) {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] stability probe missed but transport remains alive; preserving connection remote_device_id={} local_device_id={} local_node_id={} remote_node_id={} connection_id={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str,
connection_id
)));
}
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
let pending_record = self
.connection_manager
.get_by_connection_id(&connection_id)
.await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] upsert_pending remote_device_id={} remote_node_id={} connection_id={} pending_record_exists={} record={:?}",
remote_device_id,
node_id_str,
connection_id,
pending_record.is_some(),
pending_record,
)));
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(node_id_str.clone()),
transport_stable_id,
Some(connect_mode.to_string()),
)
.await;
let _ = self
.report_transport_status_for_current_generation(
&connection_id,
"iroh-relay",
None,
)
.await;
if let Some(ref token) = extracted_token {
match self
.present_and_accept_session_token(
endpoint_id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
Some(remote_device_id.clone()),
)
.await
{
Ok(approval_scope) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] session-token presented and approved connection_id={} scope={} device_id={} token_fp={} local_admission=accepted",
connection_id,
approval_scope,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str())
)));
}
Err(error) => {
if session_token_presentation_failure_action(&error)
== SessionTokenPresentationFailureAction::RetryOnNextTransport
{
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] transient session-token presentation failure; preserving peer session for canonical transport retry connection_id={} device_id={} token_fp={} error={}",
connection_id,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str()),
error
)));
return;
}
let rejection_reason =
format!("session-token-presentation-failed: {}", error);
self.reject_session_connection(
&connection_id,
rejection_reason.as_str(),
);
self.connection_manager
.set_failed(&connection_id, Some(rejection_reason))
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_SESSION_ADMISSION_REJECTED,
)
.await;
self.forget_session_connection(&connection_id);
register_auto_connect_admission_rejection(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] session-token presentation failed connection_id={} device_id={} token_fp={} error={}",
connection_id,
remote_device_id,
super::core_impl::log_fingerprint(token.as_str()),
error
)));
return;
}
}
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
let connected_record = self
.connection_manager
.get_by_connection_id(&connection_id)
.await;
let record_count = self.connection_manager.list_all().await.len();
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][wasm] set_connected remote_device_id={} remote_node_id={} connection_id={} connected_record_exists={} total_records={} record={:?}",
remote_device_id,
node_id_str,
connection_id,
connected_record.is_some(),
record_count,
connected_record,
)));
self.emit_current_wasm_connection_state(&connection_id)
.await;
clear_auto_connect_failure_state(
&remote_device_id,
failure_count,
failure_backoff_until,
);
non_initiator_escalation_count.remove(&remote_device_id);
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] connected remote_device_id={} local_device_id={} local_node_id={} remote_node_id={}",
remote_device_id,
local_device_id,
local_node_id,
node_id_str
)));
}
Err(e) => {
let err = e.to_string();
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(node_id_str.clone()),
Some(remote_device_id.clone()),
Some(node_id_str.clone()),
)
.await;
self.connection_manager
.set_failed(&connection_id, Some(err.clone()))
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
register_auto_connect_failure(
&remote_device_id,
now_millis_i64(),
failure_count,
failure_backoff_until,
);
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc] auto-connect failed to {}: {}",
remote_device_id, err
)));
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn run_auto_connect_snapshot(
self: &Arc<Self>,
user_id: &str,
local_device_id: &str,
last_attempt_at: &mut std::collections::HashMap<String, i64>,
failure_count: &mut std::collections::HashMap<String, u8>,
failure_backoff_until: &mut std::collections::HashMap<String, i64>,
retry_throttle_ms: i64,
non_initiator_wait_started_at: &mut std::collections::HashMap<String, i64>,
non_initiator_escalation_count: &mut std::collections::HashMap<String, u8>,
initiator_grace_ms: i64,
skip_log_at: &mut std::collections::HashMap<String, i64>,
connected_health_checked_at: &mut std::collections::HashMap<String, i64>,
connected_health_state: &mut std::collections::HashMap<String, bool>,
connected_health_failures: &mut std::collections::HashMap<String, u8>,
connected_health_probe_ms: i64,
connected_health_failure_threshold: u8,
last_presence_republish_at_ms: &mut i64,
presence_republish_interval_ms: i64,
last_network_change_at_ms: &mut i64,
network_change_recovery_interval_ms: i64,
network_change_failure_threshold: u8,
last_snapshot_signature: &mut Option<(usize, usize, usize)>,
last_snapshot_log_at: &mut i64,
last_known_node_id: &mut std::collections::HashMap<String, String>,
known_devices: &mut std::collections::HashMap<String, crate::signaling::Device>,
) {
if self.is_app_backgrounded() {
return;
}
let now = now_millis_i64();
match self.search_devices(user_id).await {
Ok(devices) => {
let total = devices.len();
let online = devices.iter().filter(|d| d.online).count();
let remote_online = devices
.iter()
.filter(|d| d.online && d.device_id != local_device_id)
.count();
let signature = (total, online, remote_online);
let changed = last_snapshot_signature
.map(|prev| prev != signature)
.unwrap_or(true);
let should_log = if auto_connect_verbose() {
now.saturating_sub(*last_snapshot_log_at) >= 2_000
} else {
changed || now.saturating_sub(*last_snapshot_log_at) >= 30_000
};
if should_log {
*last_snapshot_log_at = now;
*last_snapshot_signature = Some(signature);
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect][snapshot] user_id={} local_device_id={} total_devices={} online_devices={} remote_online_devices={}",
user_id, local_device_id, total, online, remote_online
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][snapshot] user_id={} local_device_id={} total_devices={} online_devices={} remote_online_devices={}",
user_id, local_device_id, total, online, remote_online
)));
}
let known_device_ids: std::collections::HashSet<String> =
devices.iter().map(|d| d.device_id.clone()).collect();
for device in devices {
known_devices.insert(device.device_id.clone(), device.clone());
self.try_auto_connect_device(
user_id,
device,
local_device_id,
last_attempt_at,
failure_count,
failure_backoff_until,
retry_throttle_ms,
non_initiator_wait_started_at,
non_initiator_escalation_count,
initiator_grace_ms,
skip_log_at,
connected_health_checked_at,
connected_health_state,
connected_health_failures,
connected_health_probe_ms,
connected_health_failure_threshold,
last_presence_republish_at_ms,
presence_republish_interval_ms,
last_network_change_at_ms,
network_change_recovery_interval_ms,
network_change_failure_threshold,
last_known_node_id,
)
.await;
}
last_attempt_at.retain(|k, _| known_device_ids.contains(k));
failure_count.retain(|k, _| known_device_ids.contains(k));
failure_backoff_until.retain(|k, _| known_device_ids.contains(k));
non_initiator_wait_started_at.retain(|k, _| known_device_ids.contains(k));
non_initiator_escalation_count.retain(|k, _| known_device_ids.contains(k));
connected_health_checked_at.retain(|k, _| known_device_ids.contains(k));
connected_health_state.retain(|k, _| known_device_ids.contains(k));
connected_health_failures.retain(|k, _| known_device_ids.contains(k));
last_known_node_id.retain(|key, _| {
let device_id = key
.strip_suffix("::runtime-instance")
.or_else(|| key.strip_suffix("::ticket"))
.unwrap_or(key.as_str());
known_device_ids.contains(device_id)
});
skip_log_at.retain(|k, _| {
k.split_once("::")
.map_or(false, |(id, _)| known_device_ids.contains(id))
});
}
Err(error) => {
let min_interval = if auto_connect_verbose() {
2_000
} else {
30_000
};
if now.saturating_sub(*last_snapshot_log_at) >= min_interval {
*last_snapshot_log_at = now;
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[pluto-rtc][auto-connect][snapshot] user_id={} local_device_id={} search_devices_error={}",
user_id, local_device_id, error
);
#[cfg(target_arch = "wasm32")]
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect][snapshot] user_id={} local_device_id={} search_devices_error={}",
user_id, local_device_id, error
)));
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn auto_connect_loop(
self: Arc<Self>,
user_id: String,
local_device_id: String,
generation: u64,
) {
use futures::StreamExt;
let mut last_attempt_at: std::collections::HashMap<String, i64> =
std::collections::HashMap::new();
let mut non_initiator_wait_started_at: std::collections::HashMap<String, i64> =
std::collections::HashMap::new();
let mut non_initiator_escalation_count: std::collections::HashMap<String, u8> =
std::collections::HashMap::new();
let mut skip_log_at: std::collections::HashMap<String, i64> =
std::collections::HashMap::new();
let mut connected_health_checked_at: std::collections::HashMap<String, i64> =
std::collections::HashMap::new();
let mut connected_health_state: std::collections::HashMap<String, bool> =
std::collections::HashMap::new();
let mut connected_health_failures: std::collections::HashMap<String, u8> =
std::collections::HashMap::new();
let mut failure_count: std::collections::HashMap<String, u8> =
std::collections::HashMap::new();
let mut failure_backoff_until: std::collections::HashMap<String, i64> =
std::collections::HashMap::new();
let mut last_known_node_id: std::collections::HashMap<String, String> =
std::collections::HashMap::new();
let mut known_devices: std::collections::HashMap<String, crate::signaling::Device> =
std::collections::HashMap::new();
let mut last_snapshot_signature: Option<(usize, usize, usize)> = None;
let mut last_snapshot_log_at = 0_i64;
let mut last_presence_republish_at_ms = 0_i64;
let retry_throttle_ms = 1_000;
let rescan_interval_ms_fg: u64 = 300_000;
let rescan_interval_ms_bg: u64 = 300_000;
let initiator_grace_ms = 3_000;
let connected_health_probe_ms_fg = 5_000_i64;
let connected_health_probe_ms_bg = i64::MAX;
let connected_health_failure_threshold = 2;
let presence_republish_interval_ms_fg = 10_000_i64;
let presence_republish_interval_ms_bg = i64::MAX;
let mut last_network_change_at_ms = 0_i64;
let network_change_recovery_interval_ms_fg = 15_000_i64;
let network_change_recovery_interval_ms_bg = i64::MAX;
let network_change_failure_threshold = auto_connect_network_change_failure_threshold();
loop {
if !self.is_auto_connect_generation_current(generation) {
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] stopping stale loop user_id={} local_device_id={} generation={}",
user_id, local_device_id, generation
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] stopping stale loop user_id={} local_device_id={} generation={}",
user_id, local_device_id, generation
)));
break;
}
if !self.is_app_backgrounded() {
self.run_auto_connect_snapshot(
&user_id,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms_fg,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms_fg,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms_fg,
network_change_failure_threshold,
&mut last_snapshot_signature,
&mut last_snapshot_log_at,
&mut last_known_node_id,
&mut known_devices,
)
.await;
}
#[allow(unused_mut)]
let mut stream = match self.subscribe_devices(&user_id).await {
Ok(s) => {
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] subscribed 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] subscribed user_id={} local_device_id={}",
user_id, local_device_id
)));
s
}
Err(e) => {
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[pluto-rtc][auto-connect] subscribe failed user_id={} local_device_id={} error={}",
user_id, local_device_id, e
);
#[cfg(target_arch = "wasm32")]
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] subscribe failed user_id={} local_device_id={} error={}",
user_id, local_device_id, e
)));
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(1500)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(1500)).await;
continue;
}
};
#[cfg(not(target_arch = "wasm32"))]
{
let mut last_bg_state = self.is_app_backgrounded();
let mut rescan_interval =
tokio::time::interval(std::time::Duration::from_millis(if last_bg_state {
rescan_interval_ms_bg
} else {
rescan_interval_ms_fg
}));
let mut health_interval = tokio::time::interval(std::time::Duration::from_millis(
connected_health_probe_ms_fg as u64,
));
health_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
if !self.is_auto_connect_generation_current(generation) {
break;
}
let bg = self.is_app_backgrounded();
if bg != last_bg_state {
last_bg_state = bg;
rescan_interval =
tokio::time::interval(std::time::Duration::from_millis(if bg {
rescan_interval_ms_bg
} else {
rescan_interval_ms_fg
}));
}
let connected_health_probe_ms = if bg {
connected_health_probe_ms_bg
} else {
connected_health_probe_ms_fg
};
let presence_republish_interval_ms = if bg {
presence_republish_interval_ms_bg
} else {
presence_republish_interval_ms_fg
};
let network_change_recovery_interval_ms = if bg {
network_change_recovery_interval_ms_bg
} else {
network_change_recovery_interval_ms_fg
};
let grace_wakeup_delay_ms = next_non_initiator_wakeup_delay_ms(
&non_initiator_wait_started_at,
&non_initiator_escalation_count,
now_millis_i64(),
);
let grace_wakeup = tokio::time::sleep(std::time::Duration::from_millis(
grace_wakeup_delay_ms.unwrap_or(rescan_interval_ms_fg),
));
tokio::pin!(grace_wakeup);
tokio::select! {
_ = rescan_interval.tick() => {
if self.is_app_backgrounded() {
continue;
}
self.run_auto_connect_snapshot(
&user_id,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms,
network_change_failure_threshold,
&mut last_snapshot_signature,
&mut last_snapshot_log_at,
&mut last_known_node_id,
&mut known_devices,
).await;
}
_ = health_interval.tick() => {
if self.is_app_backgrounded() {
continue;
}
let cached_devices: Vec<crate::signaling::Device> =
known_devices.values().cloned().collect();
for mut device in cached_devices {
if !device.online {
let Some(ticket) = device.ticket.as_deref() else {
continue;
};
let (iroh_ticket, _) =
crate::session_token::split_compound_ticket(ticket.trim());
let Ok(endpoint_addr) = parse_endpoint_ticket(iroh_ticket) else {
continue;
};
if !self.is_connected(endpoint_addr.id).await {
continue;
}
device.online = true;
}
self.try_auto_connect_device(
&user_id,
device,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms_fg,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms_fg,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms_fg,
network_change_failure_threshold,
&mut last_known_node_id,
)
.await;
}
}
_ = &mut grace_wakeup, if grace_wakeup_delay_ms.is_some() => {
if self.is_app_backgrounded() {
continue;
}
self.run_auto_connect_snapshot(
&user_id,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms,
network_change_failure_threshold,
&mut last_snapshot_signature,
&mut last_snapshot_log_at,
&mut last_known_node_id,
&mut known_devices,
).await;
}
result = stream.next() => {
let Some(result) = result else { break; };
let events = match result {
Ok(e) => e,
Err(error) => {
eprintln!(
"[pluto-rtc][auto-connect] discovery subscription ended; recreating owner-managed stream: {}",
error
);
break;
}
};
for event in events {
match event {
crate::signaling::DeviceEvent::Added { device }
| crate::signaling::DeviceEvent::Modified { device } => {
known_devices
.insert(device.device_id.clone(), device.clone());
self.try_auto_connect_device(
&user_id,
device,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms_fg,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms_fg,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms_fg,
network_change_failure_threshold,
&mut last_known_node_id,
)
.await;
}
crate::signaling::DeviceEvent::Removed { device_id } => {
known_devices.remove(&device_id);
last_attempt_at.remove(&device_id);
non_initiator_wait_started_at.remove(&device_id);
non_initiator_escalation_count.remove(&device_id);
skip_log_at.retain(|key, _| !key.starts_with(&format!("{}::", device_id)));
connected_health_checked_at.remove(&device_id);
connected_health_state.remove(&device_id);
connected_health_failures.remove(&device_id);
last_known_node_id
.remove(&format!("{device_id}::runtime-instance"));
last_known_node_id
.remove(&format!("{device_id}::ticket"));
clear_auto_connect_failure_state(
&device_id,
&mut failure_count,
&mut failure_backoff_until,
);
for record in self.connection_manager.get_by_device_id(&device_id).await {
if matches!(
record.state,
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting
) {
self.connection_manager
.set_failed(
&record.connection_id,
Some("device removed from signaling snapshot".to_string()),
)
.await;
}
}
}
}
}
}
}
}
}
#[cfg(target_arch = "wasm32")]
{
let mut stream = stream.fuse();
loop {
if !self.is_auto_connect_generation_current(generation) {
break;
}
let bg = self.is_app_backgrounded();
let connected_health_probe_ms = if bg {
connected_health_probe_ms_bg
} else {
connected_health_probe_ms_fg
};
let presence_republish_interval_ms = if bg {
presence_republish_interval_ms_bg
} else {
presence_republish_interval_ms_fg
};
let network_change_recovery_interval_ms = if bg {
network_change_recovery_interval_ms_bg
} else {
network_change_recovery_interval_ms_fg
};
let rescan_ms = if bg {
rescan_interval_ms_bg
} else {
rescan_interval_ms_fg
};
let sleep =
gloo_timers::future::sleep(std::time::Duration::from_millis(rescan_ms))
.fuse();
futures::pin_mut!(sleep);
futures::select! {
_ = sleep => {
if self.is_app_backgrounded() {
continue;
}
self.run_auto_connect_snapshot(
&user_id,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms,
network_change_failure_threshold,
&mut last_snapshot_signature,
&mut last_snapshot_log_at,
&mut last_known_node_id,
&mut known_devices,
).await;
}
result = stream.next() => {
let Some(result) = result else { break; };
let events = match result {
Ok(e) => e,
Err(error) => {
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(
&format!(
"[pluto-rtc][auto-connect] discovery subscription ended; recreating owner-managed stream: {}",
error
),
));
break;
}
};
for event in events {
match event {
crate::signaling::DeviceEvent::Added { device }
| crate::signaling::DeviceEvent::Modified { device } => {
self.try_auto_connect_device(
&user_id,
device,
&local_device_id,
&mut last_attempt_at,
&mut failure_count,
&mut failure_backoff_until,
retry_throttle_ms,
&mut non_initiator_wait_started_at,
&mut non_initiator_escalation_count,
initiator_grace_ms,
&mut skip_log_at,
&mut connected_health_checked_at,
&mut connected_health_state,
&mut connected_health_failures,
connected_health_probe_ms_fg,
connected_health_failure_threshold,
&mut last_presence_republish_at_ms,
presence_republish_interval_ms_fg,
&mut last_network_change_at_ms,
network_change_recovery_interval_ms_fg,
network_change_failure_threshold,
&mut last_known_node_id,
)
.await;
}
crate::signaling::DeviceEvent::Removed { device_id } => {
last_attempt_at.remove(&device_id);
non_initiator_wait_started_at.remove(&device_id);
non_initiator_escalation_count.remove(&device_id);
skip_log_at.retain(|key, _| !key.starts_with(&format!("{}::", device_id)));
connected_health_checked_at.remove(&device_id);
connected_health_state.remove(&device_id);
connected_health_failures.remove(&device_id);
last_known_node_id
.remove(&format!("{device_id}::runtime-instance"));
last_known_node_id
.remove(&format!("{device_id}::ticket"));
clear_auto_connect_failure_state(
&device_id,
&mut failure_count,
&mut failure_backoff_until,
);
for record in self.connection_manager.get_by_device_id(&device_id).await {
if matches!(
record.state,
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting
) {
self.connection_manager
.set_failed(
&record.connection_id,
Some("device removed from signaling snapshot".to_string()),
)
.await;
}
}
}
}
}
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
println!(
"[pluto-rtc][auto-connect] device stream ended user_id={} local_device_id={}, resubscribing",
user_id, local_device_id
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][auto-connect] device stream ended user_id={} local_device_id={}, resubscribing",
user_id, local_device_id
)));
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(400)).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(std::time::Duration::from_millis(400)).await;
}
}
}