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,
) -> 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;
let request_reciprocal = !has_current_inbound_admission_proof && !peer_requested_reciprocal;
SessionAdmissionPresentationDecision {
should_present: !has_current_remote_admission_proof
|| request_reciprocal
|| (peer_requested_reciprocal && !bilateral_proof_is_current),
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"))]
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,
) -> ActiveConnectionProbeDecision {
if independent_route_ready {
ActiveConnectionProbeDecision::Healthy
} else {
active_connection_probe_decision(active_probe_healthy, passive_transport_alive)
}
}
#[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 retryable_session_admission_failure_decision(
failure_count: u8,
failure_threshold: u8,
attempted_transport_stable_id: Option<u64>,
current_transport_stable_id: Option<u64>,
) -> RetryableSessionAdmissionFailureDecision {
if attempted_transport_stable_id != current_transport_stable_id
|| attempted_transport_stable_id.is_none()
{
return RetryableSessionAdmissionFailureDecision::ReplacementAlreadyWon;
}
if 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]")
{
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 {
if revision <= self.revision {
return false;
}
self.revision = revision;
self.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();
true
}
}
#[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(),
)
}
#[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,
}
}
#[cfg(target_arch = "wasm32")]
#[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)
&& actor.mailbox.revision == revision
&& actor
.mailbox
.peers
.get(&peer.device_id)
.is_some_and(|current| browser_peer_evidence(current) == browser_peer_evidence(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) {
client
.retire_remote_excluded_connections(&peer.device_id, binding_node_id)
.await;
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| connection.stable_id() as u64);
if !actor
.try_borrow()
.is_ok_and(|state| browser_attempt_is_current(&state, revision, peer))
{
retire_stale_browser_attempt(client, endpoint_id, transport_stable_id).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));
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
BrowserPeerReconcileOutcome::WakeAt(deadline) => {
state
.retry_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
BrowserPeerReconcileOutcome::WaitUntil(deadline) => {
state
.non_initiator_not_before_ms
.insert(peer.device_id.clone(), deadline);
}
}
}
}
#[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) = {
let mut state = actor.borrow_mut();
if !browser_actor_is_current(&state) || !state.mailbox.accept(revision, peers) {
(false, Vec::new())
} else {
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)
}
};
if accepted {
for (device_id, node_id) in current_bindings {
self.observe_authoritative_device_node(&device_id, &node_id);
}
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,
));
}
#[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,
));
}
#[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,
));
}
#[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(2, vec![peer]));
assert!(!mailbox.accept(1, Vec::new()));
assert!(!mailbox.accept(2, Vec::new()));
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"]);
assert!(mailbox.accept(3, Vec::new()));
assert!(mailbox.peers.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 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);
}
#[tokio::test]
async fn remote_exclusion_retires_existing_managed_connection() {
let client = Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
"test-app".to_string(),
Box::new(|| None),
);
client
.connection_manager
.upsert_pending(
"conn-remote-excluded".to_string(),
Some("not-an-endpoint".to_string()),
Some("device-remote".to_string()),
Some("not-an-endpoint".to_string()),
)
.await;
client
.connection_manager
.set_connected("conn-remote-excluded", Some("not-an-endpoint".to_string()))
.await;
client
.retire_remote_excluded_connections("device-remote", Some("not-an-endpoint"))
.await;
assert!(client
.connection_manager
.get_by_connection_id("conn-remote-excluded")
.await
.is_none());
}
}
fn should_suppress_auto_connect_for_replacement_window(
snapshot: &ManagedConnectionHealthSnapshot,
now_ms: i64,
) -> bool {
if !snapshot.replacement_pending {
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_remote_excluded_connections(
&self,
remote_device_id: &str,
remote_node_id: Option<&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 reason = crate::lifecycle_reason::REASON_REMOTE_MANUAL_DISCONNECT;
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 remote-excluded connection remote_device_id={} connection_id={} transport_closed={}",
remote_device_id, record.connection_id, disconnected
);
}
}
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) {
self.retire_remote_excluded_connections(&remote_device_id, device.node_id.as_deref())
.await;
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"))]
{
let mut credentials = self
.native_route_repair_credentials
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let (Some(token), Some(token_payload)) =
(extracted_token.as_ref(), extracted_token_suffix.as_ref())
{
credentials.insert(
node_id_str.clone(),
crate::client::NativeRouteRepairCredential {
token: token.clone(),
token_payload: token_payload.clone(),
authoritative_device_id: remote_device_id.clone(),
},
);
} else {
credentials.remove(&node_id_str);
}
}
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| connection.stable_id() as u64);
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,
) {
println!(
"[pluto-rtc][auto-connect] peer runtime instance changed remote_device_id={} node_id={} - replacing unresponsive physical transport",
remote_device_id, node_id_str
);
if let Some(transport_stable_id) = current_transport_stable_id {
let _ = self
.disconnect_transport_generation_with_reason(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.await;
}
} 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;
if let Some(snapshot) = managed_health_snapshot.as_ref() {
if should_suppress_auto_connect_for_replacement_window(&snapshot, now) {
log_skip("replacement-pending");
return;
}
}
#[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;
#[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 probe_decision = active_connection_probe_decision_with_independent_route(
active_probe_healthy,
passive_transport_alive,
independent_route_ready,
);
healthy = !matches!(probe_decision, ActiveConnectionProbeDecision::Failed);
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 if matches!(probe_decision, ActiveConnectionProbeDecision::Failed) {
let failures = connected_health_failures
.entry(remote_device_id.clone())
.or_insert(0);
*failures = failures.saturating_add(1).min(10);
} else {
connected_health_failures.remove(&remote_device_id);
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;
}
log_skip("stale-active-connection-reconnect");
self.maybe_republish_presence_for_auto_connect(
user_id,
local_device_id,
"stale-active-connection-reconnect",
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
if let Some(transport_stable_id) = probed_transport_stable_id {
let _ = self
.disconnect_transport_generation_with_reason(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT,
)
.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 {
connected_health_failures.remove(&remote_device_id);
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(|conn| conn.stable_id() as u64);
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,
);
#[cfg(target_arch = "wasm32")]
let presentation_decision = session_admission_presentation_decision(
extracted_token.is_some(),
has_remote_admission_proof,
true,
false,
);
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()),
has_remote_admission_proof,
presentation_decision.request_reciprocal,
)
.await;
#[cfg(target_arch = "wasm32")]
let presentation_result = self
.present_and_accept_session_token(
endpoint_id,
&connection_id,
token,
extracted_token_suffix.as_deref(),
Some(remote_device_id.clone()),
)
.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| connection.stable_id() as u64);
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,
) {
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");
}
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(node_id_str.clone()),
self.get_connection(endpoint_id)
.await
.map(|conn| conn.stable_id() as u64),
Some(connect_mode.to_string()),
)
.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) => {
#[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_republish_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| connection.stable_id() as u64);
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>,
) {
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 {
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 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 = 5_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,
)
.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
}));
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
};
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,
).await;
}
result = stream.next() => {
let Some(result) = result else { break; };
let events = match result {
Ok(e) => e,
Err(_) => continue,
};
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(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,
).await;
}
result = stream.next() => {
let Some(result) = result else { break; };
let events = match result {
Ok(e) => e,
Err(_) => continue,
};
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;
}
}
}