use super::*;
#[cfg(not(target_arch = "wasm32"))]
use iroh::Watcher;
pub(crate) fn dial_timeout_reason(label: &str) -> String {
format!(
"{}:{label}-timeout",
crate::lifecycle_reason::REASON_AUTO_CONNECT_FAILURE
)
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
fn install_compiled_packet_carriers(
mut builder: iroh::endpoint::Builder,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(
iroh::endpoint::Builder,
HashMap<
crate::iroh_carrier_kind::IrohCarrierKind,
Arc<crate::packet_carrier_transport::PacketCarrierAddressProvider>,
>,
Arc<crate::packet_carrier_transport::PacketCarrierPathSelector>,
)> {
use crate::iroh_carrier_kind::IrohCarrierKind;
#[cfg(feature = "transport-moq")]
use crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID;
#[cfg(feature = "transport-webrtc")]
use crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID;
let mut carriers = HashMap::new();
let path_selector =
Arc::new(crate::packet_carrier_transport::PacketCarrierPathSelector::default());
builder = builder.path_selector(path_selector.clone());
#[cfg(feature = "transport-webrtc")]
{
let transport = Arc::new(
crate::packet_carrier_transport::PacketCarrierTransport::new(
EXPERIMENTAL_WEBRTC_TRANSPORT_ID,
endpoint_id,
)?,
);
builder = builder.add_custom_transport(transport.clone());
carriers.insert(
IrohCarrierKind::WebRtc,
Arc::new(
crate::packet_carrier_transport::PacketCarrierAddressProvider::new(
IrohCarrierKind::WebRtc,
transport,
),
),
);
}
#[cfg(feature = "transport-moq")]
{
let transport = Arc::new(
crate::packet_carrier_transport::PacketCarrierTransport::new(
EXPERIMENTAL_MOQ_TRANSPORT_ID,
endpoint_id,
)?,
);
builder = builder.add_custom_transport(transport.clone());
carriers.insert(
IrohCarrierKind::Moq,
Arc::new(
crate::packet_carrier_transport::PacketCarrierAddressProvider::new(
IrohCarrierKind::Moq,
transport,
),
),
);
}
Ok((builder, carriers, path_selector))
}
fn parse_loopback_test_relay_url(raw_url: Option<&str>) -> anyhow::Result<Option<iroh::RelayUrl>> {
let Some(raw_url) = raw_url.map(str::trim).filter(|value| !value.is_empty()) else {
return Ok(None);
};
let relay_url: iroh::RelayUrl = raw_url
.parse()
.map_err(|error| anyhow::anyhow!("invalid local Iroh relay URL: {error}"))?;
let host = relay_url
.host_str()
.unwrap_or_default()
.trim_end_matches('.');
let is_loopback =
host.eq_ignore_ascii_case("localhost") || host == "127.0.0.1" || host == "::1";
anyhow::ensure!(
relay_url.scheme() == "https" && is_loopback,
"local Iroh relay URL must use HTTPS and a loopback host"
);
anyhow::ensure!(
relay_url.query().is_none() && relay_url.fragment().is_none(),
"local Iroh relay URL must not contain a query or fragment"
);
Ok(Some(relay_url))
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn is_duplicate_kept_existing_close_reason(error: Option<&str>) -> bool {
matches!(
crate::lifecycle_reason::LifecycleReasonCode::from_text(error),
Some(crate::lifecycle_reason::LifecycleReasonCode::DuplicateKeptExisting)
)
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn is_replacement_churn_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::is_replacement_churn(error)
}
#[cfg(test)]
pub(crate) fn is_graceful_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::is_graceful(error)
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn is_terminal_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::is_terminal(error)
}
pub(crate) fn is_manual_disconnect_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::is_manual_disconnect(error)
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn session_admission_timeout_owns_current_transport(
admission_pending: bool,
admission_response_in_flight: bool,
timeout_transport_stable_id: u64,
active_transport_stable_id: Option<u64>,
managed_transport_stable_id: Option<u64>,
) -> bool {
admission_pending
&& !admission_response_in_flight
&& active_transport_stable_id == Some(timeout_transport_stable_id)
&& managed_transport_stable_id == Some(timeout_transport_stable_id)
}
pub(crate) fn accepted_transport_event_is_current(
accepted_transport_stable_id: u64,
current_transport_stable_id: Option<u64>,
) -> bool {
current_transport_stable_id == Some(accepted_transport_stable_id)
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn desired_user_device_reopen_candidate(
needs_reopen: bool,
desired_device_id: Option<&str>,
) -> Option<&str> {
needs_reopen.then_some(desired_device_id).flatten()
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn transport_replacement_requires_application_crypto_reconfirmation(
previous_transport_stable_id: Option<u64>,
current_transport_stable_id: u64,
) -> bool {
previous_transport_stable_id.is_some_and(|previous| previous != current_transport_stable_id)
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn accepted_transport_replaces_managed_generation(
physically_replaced_transport_stable_id: Option<u64>,
accepted_transport_stable_id: u64,
projected_transport_generation: u64,
logical_session_was_admitted: bool,
) -> bool {
physically_replaced_transport_stable_id
.is_some_and(|prior| prior != accepted_transport_stable_id)
|| projected_transport_generation > 1
|| logical_session_was_admitted
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn incoming_transport_requires_carrier_restart(
already_managed_connected: bool,
replaced_main_route: bool,
) -> bool {
!already_managed_connected || replaced_main_route
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn route_repair_needed(promoted_user_device: bool, replaced_main_route: bool) -> bool {
promoted_user_device && replaced_main_route
}
#[cfg_attr(not(any(test, target_arch = "wasm32")), allow(dead_code))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum WasmClosedTransportRebind {
TerminalDisconnect,
AlreadyPromoted(u64),
BaseReplacement(u64),
SurvivingCarrierAwaitingReplacement,
NoLiveTransport,
}
#[cfg_attr(not(any(test, target_arch = "wasm32")), allow(dead_code))]
pub(crate) fn classify_wasm_closed_transport_rebind(
terminal_disconnect: bool,
current_base_stable_id: Option<u64>,
connected_managed_stable_id: Option<u64>,
has_live_transport: bool,
) -> WasmClosedTransportRebind {
if terminal_disconnect {
return WasmClosedTransportRebind::TerminalDisconnect;
}
match current_base_stable_id {
Some(stable_id) if connected_managed_stable_id == Some(stable_id) => {
WasmClosedTransportRebind::AlreadyPromoted(stable_id)
}
Some(stable_id) => WasmClosedTransportRebind::BaseReplacement(stable_id),
None if has_live_transport => {
WasmClosedTransportRebind::SurvivingCarrierAwaitingReplacement
}
None => WasmClosedTransportRebind::NoLiveTransport,
}
}
pub(crate) fn log_fingerprint(value: &str) -> String {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
value.hash(&mut hasher);
format!("fp-{:08x}", (hasher.finish() & 0xffff_ffff) as u32)
}
pub(crate) fn summarize_compound_ticket_for_logs(
ticket: &str,
) -> (String, Option<String>, Option<String>) {
let (iroh_ticket, suffix) = crate::session_token::split_ticket(ticket);
let iroh_fingerprint = log_fingerprint(iroh_ticket);
let payload = suffix.and_then(crate::session_token::decode_token_payload);
let scope = payload.as_ref().map(|value| value.scope.to_string());
let token_fingerprint = payload
.as_ref()
.map(|value| log_fingerprint(value.token.as_str()));
(iroh_fingerprint, scope, token_fingerprint)
}
#[cfg(not(target_arch = "wasm32"))]
pub(super) const PERSISTENT_MANAGED_ADMISSION_SCOPE: &str = "user-device";
#[cfg(not(target_arch = "wasm32"))]
pub(super) const PERSISTENT_MANAGED_SCOPE_TICKET_FILE_PREFIX: &str = "openrtc_managed_scope_ticket";
#[cfg(not(target_arch = "wasm32"))]
const RELAY_ONLY_ENDPOINT_TICKET_RETRY_WINDOW: std::time::Duration =
std::time::Duration::from_secs(12);
#[cfg(not(target_arch = "wasm32"))]
const RELAY_ONLY_ENDPOINT_TICKET_MIN_RETRY_DELAY: std::time::Duration =
std::time::Duration::from_millis(250);
#[cfg(not(target_arch = "wasm32"))]
const RELAY_ONLY_ENDPOINT_TICKET_MAX_RETRY_DELAY: std::time::Duration =
std::time::Duration::from_millis(1_500);
impl Client {
#[cfg(not(target_arch = "wasm32"))]
pub fn offline(&self) -> crate::offline::OfflineClient<'_> {
crate::offline::OfflineClient::new(self)
}
pub fn new(api_key: String) -> anyhow::Result<Self> {
let api_key = crate::validate_api_key(&api_key)?;
Ok(ClientBuilder::new_provider_neutral(
crate::app_tag_from_api_key(api_key),
Box::new(|| None),
)
.build())
}
#[cfg(not(target_arch = "wasm32"))]
async fn admission_timeout_still_owns_transport(
&self,
connection_id: &str,
remote_endpoint_id: iroh::EndpointId,
timeout_transport_stable_id: u64,
) -> bool {
let admission_pending = matches!(
self.session_admission(connection_id),
crate::session_token::SessionAdmission::Pending
);
let admission_response_in_flight = self
.session_token_registry
.is_session_token_admission_in_flight(connection_id);
let active_transport_stable_id = self
.get_connection(remote_endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let managed_transport_stable_id = self
.connection_manager
.get_by_connection_id(connection_id)
.await
.and_then(|record| record.transport_stable_id);
let owns = session_admission_timeout_owns_current_transport(
admission_pending,
admission_response_in_flight,
timeout_transport_stable_id,
active_transport_stable_id,
managed_transport_stable_id,
);
if !owns {
eprintln!(
"[PlutoRTC][session-admission][timeout-fenced] connection_id={} remote_endpoint_id={} timeout_stable_id={} admission_pending={} response_in_flight={} active_stable_id={:?} managed_stable_id={:?}",
connection_id,
remote_endpoint_id,
timeout_transport_stable_id,
admission_pending,
admission_response_in_flight,
active_transport_stable_id,
managed_transport_stable_id,
);
}
owns
}
#[cfg_attr(not(any(test, target_arch = "wasm32")), allow(dead_code))]
#[cfg(test)]
#[deprecated(
since = "2.0.0",
note = "rollback-only: use Client::new or Client::builder"
)]
pub fn new_with_app_tag(
project_id: String,
app_tag: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
ClientBuilder::new_with_app_tag(project_id, app_tag, token_provider).build()
}
pub fn new_provider_neutral(
app_tag: String,
identity_credential_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
ClientBuilder::new_provider_neutral(app_tag, identity_credential_provider).build()
}
#[cfg(test)]
pub(crate) fn new_for_test(
project_id: String,
api_key: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
ClientBuilder::new_for_test(project_id, api_key, token_provider).build()
}
pub fn builder_provider_neutral(
app_tag: String,
identity_credential_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> ClientBuilder {
ClientBuilder::new_provider_neutral(app_tag, identity_credential_provider)
}
pub fn builder(
api_key: String,
identity_credential_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> anyhow::Result<ClientBuilder> {
let api_key = crate::validate_api_key(&api_key)?;
Ok(ClientBuilder::new_provider_neutral(
crate::app_tag_from_api_key(api_key),
identity_credential_provider,
))
}
pub fn app_tag(&self) -> &str {
&self.app_tag
}
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
pub async fn enable_scoped_connection_actor(
&self,
registry: Arc<crate::client::scoped_connection_actor::ConnectionActors>,
) {
let mut slot = self.scoped_connection_actor_registry.write().await;
if let Some(prev) = slot.take() {
prev.shutdown_all().await;
}
*slot = Some(registry);
}
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
pub async fn scoped_connection_actor_registry(
&self,
) -> Option<Arc<crate::client::scoped_connection_actor::ConnectionActors>> {
self.scoped_connection_actor_registry.read().await.clone()
}
#[cfg(all(not(target_arch = "wasm32"), feature = "experimental-scoped-actor"))]
pub async fn disable_scoped_connection_actor(&self) {
let mut slot = self.scoped_connection_actor_registry.write().await;
if let Some(prev) = slot.take() {
prev.shutdown_all().await;
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn auth_readiness(&self) -> Arc<crate::client::auth_readiness::AuthReadinessStore> {
self.auth_readiness.clone()
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn correlation_for_connection(
&self,
connection_id: &str,
) -> crate::client::correlation::CorrelationContext {
use crate::client::correlation::CorrelationContext;
let mut ctx = CorrelationContext::new().connection_id(connection_id);
if let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
{
if let Some(node_id) = record.node_id.as_deref().or(record.endpoint_id.as_deref()) {
ctx = ctx.peer_node_id(node_id);
}
}
let scopes = self.connection_manager.get_scopes(connection_id).await;
let classifier = self.scope_classifier.read().await.clone();
let classification = classifier.classify_scopes(&scopes);
if let Some(scope) = classification.scope.as_deref() {
ctx = ctx.scope(scope);
}
if let Some(grant_id) = classification.grant_id.as_deref() {
ctx = ctx.grant_id(grant_id);
}
if let Some(kind) = classification.session_kind {
ctx = ctx.session_kind(kind);
}
ctx
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn init_native_device_identity(
&self,
base_dir: std::path::PathBuf,
preferred_name: Option<&str>,
) -> anyhow::Result<crate::native_device::NativeDeviceIdentity> {
let _init_guard = self.native_device_identity_init_guard.lock().await;
let identity = crate::native_device::load_or_create(&base_dir, preferred_name).await?;
{
let mut base_dir_guard = self.native_device_base_dir.write().await;
*base_dir_guard = Some(base_dir);
}
{
let mut identity_guard = self.native_device_identity.write().await;
*identity_guard = Some(identity.clone());
}
self.rehydrate_persistent_managed_scope_ticket(PERSISTENT_MANAGED_ADMISSION_SCOPE)
.await?;
let _ = self.native_device_updates.send(identity.clone());
Ok(identity)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn get_native_device_identity(
&self,
) -> anyhow::Result<crate::native_device::NativeDeviceIdentity> {
if let Some(identity) = self.native_device_identity.read().await.clone() {
return Ok(identity);
}
let base_dir = self
.native_device_base_dir
.read()
.await
.clone()
.ok_or_else(|| anyhow::anyhow!("native device identity not initialized"))?;
self.init_native_device_identity(base_dir, None).await
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn get_native_system_device_info(
&self,
) -> anyhow::Result<crate::native_device::SystemDeviceInfo> {
let identity = self.get_native_device_identity().await?;
Ok(identity
.system_info
.unwrap_or_else(crate::native_device::system_device_info))
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn update_native_device_name(
&self,
device_name: &str,
) -> anyhow::Result<crate::native_device::NativeDeviceIdentity> {
let trimmed = device_name.trim();
if trimmed.is_empty() {
return Err(anyhow::anyhow!("device name cannot be empty"));
}
let base_dir = self
.native_device_base_dir
.read()
.await
.clone()
.ok_or_else(|| anyhow::anyhow!("native device identity not initialized"))?;
let mut identity = self.get_native_device_identity().await?;
if identity.device_name == trimmed {
return Ok(identity);
}
identity.device_name = trimmed.to_string();
identity.name_source = Some(crate::native_device::DeviceNameSource::UserProvided);
identity.updated_at_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as i64)
.unwrap_or(identity.updated_at_ms);
crate::native_device::persist(&base_dir, &identity).await?;
{
let mut identity_guard = self.native_device_identity.write().await;
*identity_guard = Some(identity.clone());
}
let _ = self.native_device_updates.send(identity.clone());
Ok(identity)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn subscribe_native_device_updates(
&self,
) -> tokio::sync::broadcast::Receiver<crate::native_device::NativeDeviceIdentity> {
self.native_device_updates.subscribe()
}
#[cfg(not(target_arch = "wasm32"))]
pub fn connection_state_updates(
&self,
) -> tokio::sync::broadcast::Receiver<crate::client::StateSnapshot> {
self.native_connection_state_updates.subscribe()
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn emit_current_native_connection_state(&self, connection_id: &str) {
if let Some(snapshot) = self.connection_state(connection_id).await {
let _ = self.native_connection_state_updates.send(snapshot);
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn metadata_with_local_device_id(
&self,
metadata: Option<&str>,
) -> Option<String> {
let device_id = self
.native_device_identity
.read()
.await
.as_ref()
.map(|identity| identity.device_id.clone());
let Some(device_id) = device_id else {
return metadata.map(str::to_string);
};
let Some(raw_metadata) = metadata else {
return Some(serde_json::json!({ "deviceId": device_id }).to_string());
};
match serde_json::from_str::<serde_json::Value>(raw_metadata) {
Ok(serde_json::Value::Object(mut map)) => {
map.entry("deviceId".to_string())
.or_insert_with(|| serde_json::Value::String(device_id));
Some(serde_json::Value::Object(map).to_string())
}
_ => Some(
serde_json::json!({
"deviceId": device_id,
"metadata": raw_metadata,
})
.to_string(),
),
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn init_iroh_with_router_mode(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
spawn_internal_router: bool,
test_relay_url: Option<&str>,
) -> anyhow::Result<String> {
let _init_lock = self.iroh_init_guard.lock().await;
{
let endpoint_guard = self.iroh_endpoint.read().await;
if let Some(endpoint) = endpoint_guard.as_ref() {
let endpoint_node_id = endpoint.id().to_string();
drop(endpoint_guard);
let mut node_guard = self.node_id.write().await;
if node_guard.as_deref() != Some(endpoint_node_id.as_str()) {
eprintln!(
"[pluto-rtc][native] repaired local node projection from live endpoint old_node_id={} endpoint_node_id={}",
node_guard.as_deref().unwrap_or("<none>"),
endpoint_node_id
);
*node_guard = Some(endpoint_node_id.clone());
}
return Ok(endpoint_node_id);
}
}
let endpoint_secret_key = match secret_key {
Some(key_bytes) => iroh::SecretKey::try_from(&key_bytes[..])?,
None => iroh::SecretKey::generate(),
};
let transport_config = self.transport_config.read().await.clone();
anyhow::ensure!(
transport_config.network_policy != crate::client::NetworkPolicy::LocalOnly
|| cfg!(feature = "transport-lan"),
"LocalOnly requires the transport-lan feature"
);
let mut builder =
if transport_config.network_policy == crate::client::NetworkPolicy::LocalOnly {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
} else {
iroh::Endpoint::builder(iroh::endpoint::presets::N0)
};
let alpns = crate::protocol_config::router_alpns(extra_alpns, false)?;
builder = builder.alpns(alpns);
if let Some(relay_url) = parse_loopback_test_relay_url(test_relay_url)? {
builder = builder.relay_mode(iroh::RelayMode::custom([relay_url]));
#[cfg(any(feature = "test-harness", feature = "test-relay-client"))]
{
builder = builder.ca_tls_config(iroh::tls::CaTlsConfig::insecure_skip_verify());
}
}
if !transport_config.relay {
builder = builder.relay_mode(iroh::RelayMode::Disabled);
}
let relay_only = should_use_relay_only_mode() || transport_config.iroh_relay_only;
let relay_transport_policy = transport_config
.iroh_relay_transport_policy
.unwrap_or(crate::client::RelayPolicy::WebsocketRequired);
builder = apply_native_network_preferences(builder, relay_only, relay_transport_policy)?;
builder = self.apply_optional_lan_discovery(builder).await?;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let packet_carriers;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let carrier_path_selector;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
let installed =
install_compiled_packet_carriers(builder, endpoint_secret_key.public())?;
builder = installed.0;
packet_carriers = installed.1;
carrier_path_selector = installed.2;
}
builder = builder.secret_key(endpoint_secret_key);
let endpoint = builder.bind().await?;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
let providers = packet_carriers.values().cloned().collect::<Vec<_>>();
*self.iroh_packet_carriers.write().await = packet_carriers;
for provider in providers {
self.add_upgrade_provider(provider).await?;
}
}
let node_id = endpoint.id().to_string();
self.start_local_discovery_tasks(&endpoint, &node_id).await;
let relay_was_online_before_publish = if !transport_config.relay {
true
} else {
match tokio::time::timeout(std::time::Duration::from_secs(10), endpoint.online()).await
{
Ok(()) => {
eprintln!("[pluto-rtc][native] relay online node_id={}", node_id);
true
}
Err(_) => {
eprintln!(
"[pluto-rtc][native] WARNING: relay not connected after 10s, \
tickets may lack relay URLs node_id={}",
node_id
);
false
}
}
};
#[cfg(not(target_arch = "wasm32"))]
if !relay_was_online_before_publish {
let endpoint_for_update = endpoint.clone();
let presence_tx = self.presence_loop_tx.clone();
tokio::spawn(async move {
endpoint_for_update.online().await;
if let Some(tx) = presence_tx.lock().unwrap().clone() {
let _ = tx.try_send(crate::presence::PresenceCommand::RepublishDurableNow);
}
});
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let node = IrohNativeNode::spawn_with_endpoint_and_carrier_path_selector(
endpoint.clone(),
spawn_internal_router,
carrier_path_selector,
)
.await?;
#[cfg(not(any(feature = "transport-webrtc", feature = "transport-moq")))]
let node = if spawn_internal_router {
IrohNativeNode::spawn_with_endpoint(endpoint.clone()).await?
} else {
IrohNativeNode::spawn_with_endpoint_no_router(endpoint.clone()).await?
};
let native_accept_events = node.accept_events();
let native_incoming_streams = node.incoming_streams_stream();
let mut node_guard = self.iroh_node.write().await;
*node_guard = Some(node);
drop(node_guard);
let mut n_guard = self.node_id.write().await;
*n_guard = Some(node_id.clone());
drop(n_guard);
let mut guard = self.iroh_endpoint.write().await;
*guard = Some(endpoint);
drop(guard);
if spawn_internal_router {
self.start_native_accept_bridge(native_accept_events);
} else {
self.start_native_external_outbound_close_bridge(native_accept_events);
}
self.start_native_incoming_stream_router(native_incoming_streams);
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[pluto-rtc][native] initialized endpoint node_id={} internal_router={}",
node_id, spawn_internal_router
);
Ok(node_id)
}
pub async fn init_iroh(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
) -> anyhow::Result<String> {
#[cfg(not(target_arch = "wasm32"))]
{
return self
.init_iroh_with_router_mode(secret_key, extra_alpns, true, None)
.await;
}
#[cfg(target_arch = "wasm32")]
{
wasm_init_log("init_iroh:enter");
{
let endpoint_guard = self.iroh_endpoint.read().await;
let node_guard = self.node_id.read().await;
if endpoint_guard.is_some() {
if let Some(existing_node_id) = node_guard.clone() {
wasm_init_log("init_iroh:reuse-existing-node-id");
return Ok(existing_node_id);
}
}
}
wasm_init_log("init_iroh:builder-created");
let mut builder = iroh::Endpoint::builder(iroh::endpoint::presets::N0);
let endpoint_secret_key = match secret_key {
Some(key_bytes) => {
wasm_init_log("init_iroh:using-provided-secret-key");
iroh::SecretKey::try_from(&key_bytes[..])?
}
None => iroh::SecretKey::generate(),
};
let alpns = crate::protocol_config::router_alpns(
extra_alpns,
cfg!(feature = "iroh-protocols-wasm"),
)?;
let relay = self.transport_config.read().await.relay;
builder = builder.alpns(alpns).relay_mode(if relay {
iroh::RelayMode::Default
} else {
iroh::RelayMode::Disabled
});
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let packet_carriers;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let carrier_path_selector;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
let installed =
install_compiled_packet_carriers(builder, endpoint_secret_key.public())?;
builder = installed.0;
packet_carriers = installed.1;
carrier_path_selector = installed.2;
}
builder = builder.secret_key(endpoint_secret_key);
wasm_init_log("init_iroh:bind-start");
let endpoint = builder.bind().await?;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
*self.iroh_packet_carriers.write().await = packet_carriers;
}
wasm_init_log("init_iroh:bind-complete");
let node_id = endpoint.id().to_string();
wasm_init_log("init_iroh:spawn-router-start");
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let node = IrohWasmNode::spawn_with_endpoint_and_carrier_path_selector(
endpoint.clone(),
carrier_path_selector,
)
.await?;
#[cfg(not(any(feature = "transport-webrtc", feature = "transport-moq")))]
let node = IrohWasmNode::spawn_with_endpoint(endpoint.clone()).await?;
wasm_init_log("init_iroh:spawn-router-complete");
wasm_init_log("init_iroh:node-write-start");
let mut node_guard = self.iroh_node.write().await;
*node_guard = Some(node);
drop(node_guard);
wasm_init_log("init_iroh:node-write-complete");
wasm_init_log("init_iroh:node-id-write-start");
let mut n_guard = self.node_id.write().await;
*n_guard = Some(node_id.clone());
drop(n_guard);
wasm_init_log("init_iroh:node-id-write-complete");
wasm_init_log("init_iroh:endpoint-write-start");
let mut guard = self.iroh_endpoint.write().await;
*guard = Some(endpoint);
wasm_init_log("init_iroh:endpoint-write-complete");
wasm_init_log("init_iroh:complete");
Ok(node_id)
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn init_iroh_without_internal_router(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
) -> anyhow::Result<String> {
self.init_iroh_with_router_mode(secret_key, extra_alpns, false, None)
.await
}
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "test-harness", feature = "test-relay-client")
))]
pub async fn init_test_endpoint(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
test_relay_url: Option<&str>,
) -> anyhow::Result<String> {
self.init_iroh_with_router_mode(secret_key, extra_alpns, false, test_relay_url)
.await
}
pub async fn init_iroh_with_test_relay(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
test_relay_url: Option<&str>,
) -> anyhow::Result<String> {
#[cfg(not(target_arch = "wasm32"))]
{
return self
.init_iroh_with_router_mode(secret_key, extra_alpns, true, test_relay_url)
.await;
}
#[cfg(target_arch = "wasm32")]
{
wasm_init_log("init_iroh:enter");
{
let endpoint_guard = self.iroh_endpoint.read().await;
let node_guard = self.node_id.read().await;
if endpoint_guard.is_some() {
if let Some(existing_node_id) = node_guard.clone() {
wasm_init_log("init_iroh:reuse-existing-node-id");
return Ok(existing_node_id);
}
}
}
wasm_init_log("init_iroh:builder-created");
let mut builder = iroh::Endpoint::builder(iroh::endpoint::presets::N0);
let endpoint_secret_key = match secret_key {
Some(key_bytes) => {
wasm_init_log("init_iroh:using-provided-secret-key");
iroh::SecretKey::try_from(&key_bytes[..])?
}
None => iroh::SecretKey::generate(),
};
let alpns = crate::protocol_config::router_alpns(
extra_alpns,
cfg!(feature = "iroh-protocols-wasm"),
)?;
let relay = self.transport_config.read().await.relay;
builder = builder.alpns(alpns).relay_mode(
match (relay, parse_loopback_test_relay_url(test_relay_url)?) {
(false, _) => iroh::RelayMode::Disabled,
(true, Some(relay_url)) => iroh::RelayMode::custom([relay_url]),
(true, None) => iroh::RelayMode::Default,
},
);
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let packet_carriers;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let carrier_path_selector;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
let installed =
install_compiled_packet_carriers(builder, endpoint_secret_key.public())?;
builder = installed.0;
packet_carriers = installed.1;
carrier_path_selector = installed.2;
}
builder = builder.secret_key(endpoint_secret_key);
wasm_init_log("init_iroh:bind-start");
let endpoint = builder.bind().await?;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
{
*self.iroh_packet_carriers.write().await = packet_carriers;
}
wasm_init_log("init_iroh:bind-complete");
let node_id = endpoint.id().to_string();
wasm_init_log("init_iroh:spawn-router-start");
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
let node = IrohWasmNode::spawn_with_endpoint_and_carrier_path_selector(
endpoint.clone(),
carrier_path_selector,
)
.await?;
#[cfg(not(any(feature = "transport-webrtc", feature = "transport-moq")))]
let node = IrohWasmNode::spawn_with_endpoint(endpoint.clone()).await?;
wasm_init_log("init_iroh:spawn-router-complete");
let mut node_guard = self.iroh_node.write().await;
*node_guard = Some(node);
drop(node_guard);
let mut n_guard = self.node_id.write().await;
*n_guard = Some(node_id.clone());
drop(n_guard);
let mut guard = self.iroh_endpoint.write().await;
*guard = Some(endpoint);
drop(guard);
wasm_init_log("init_iroh:complete");
Ok(node_id)
}
}
pub async fn get_endpoint(&self) -> anyhow::Result<Endpoint> {
let guard = self.iroh_endpoint.read().await;
guard
.clone()
.ok_or_else(|| anyhow::anyhow!("Iroh Endpoint not initialized"))
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
pub(crate) async fn activate_iroh_packet_carrier(
&self,
kind: crate::iroh_carrier_kind::IrohCarrierKind,
remote_endpoint_id: iroh::EndpointId,
) -> anyhow::Result<crate::packet_carrier_transport::PacketCarrierSession> {
let carriers = self.iroh_packet_carriers.read().await;
let provider = carriers
.get(&kind)
.ok_or_else(|| anyhow::anyhow!("{kind:?} packet carrier is not installed"))?;
provider
.activate(remote_endpoint_id)
.map_err(anyhow::Error::from)
}
#[cfg(all(
target_arch = "wasm32",
any(feature = "transport-webrtc", feature = "transport-moq")
))]
pub(crate) async fn active_iroh_packet_carrier_addr(
&self,
kind: crate::iroh_carrier_kind::IrohCarrierKind,
remote_endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::EndpointAddr> {
let carriers = self.iroh_packet_carriers.read().await;
let provider = carriers
.get(&kind)
.ok_or_else(|| anyhow::anyhow!("{kind:?} packet carrier is not installed"))?;
provider.prepared_endpoint_addr(remote_endpoint_id)
}
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
pub(crate) async fn close_current_iroh_carrier_generation_with_reason(
&self,
connection_id: &str,
remote_endpoint_id: iroh::EndpointId,
kind: IrohPathKind,
expected_transport_generation: u64,
reason: &str,
) -> bool {
let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return false;
};
if record.transport_generation != expected_transport_generation
|| record.active_transport != kind.transport_label()
|| self.iroh_path_kind(&remote_endpoint_id.to_string()).await != kind
{
return false;
}
let Some(connection) = self.get_connection(remote_endpoint_id).await else {
return false;
};
let stable_id = crate::transport_generation::for_connection(&connection);
if record.transport_stable_id != Some(stable_id) {
return false;
}
if crate::lifecycle_reason::is_terminal(Some(reason)) {
if self
.connection_manager
.set_closed_if_current(connection_id, stable_id, Some(reason.to_string()))
.await
.is_none()
{
return false;
}
#[cfg(target_arch = "wasm32")]
self.emit_current_wasm_connection_state(connection_id).await;
#[cfg(not(target_arch = "wasm32"))]
self.emit_current_native_connection_state(connection_id)
.await;
return self
.disconnect_with_reason(remote_endpoint_id, reason)
.await
.is_ok();
}
#[cfg(target_arch = "wasm32")]
{
let node = self.iroh_node.read().await.as_ref().cloned();
let Some(node) = node else {
return false;
};
let Ok(true) = node
.disconnect_with_reason_if_current(remote_endpoint_id, stable_id, reason)
.await
else {
return false;
};
node.clear_carrier_path_preference(remote_endpoint_id);
if self
.connection_manager
.mark_transport_replaced_if_current(
connection_id,
stable_id,
Some("wasm-carrier-generation-retired".to_string()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string()),
)
.await
.is_some()
{
self.emit_current_wasm_connection_state(connection_id).await;
}
true
}
#[cfg(not(target_arch = "wasm32"))]
{
connection.close(0u8.into(), reason.as_bytes());
true
}
}
#[cfg(all(
target_arch = "wasm32",
any(feature = "transport-webrtc", feature = "transport-moq")
))]
pub(crate) async fn recover_retired_wasm_carrier_generation(
&self,
connection_id: &str,
remote_endpoint_id: iroh::EndpointId,
expected_transport_generation: u64,
) -> bool {
let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return false;
};
let expected_endpoint_id = remote_endpoint_id.to_string();
if record.transport_generation != expected_transport_generation
|| record.transport_stable_id.is_some()
|| matches!(
record.state,
crate::connection_manager::ConnectionState::Closed
)
|| crate::lifecycle_reason::is_terminal(record.status_reason.as_deref())
|| record
.endpoint_id
.as_deref()
.or(record.node_id.as_deref())
.is_some_and(|id| id != expected_endpoint_id)
{
return false;
}
self.ensure_connected(remote_endpoint_id).await.is_ok()
}
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "transport-webrtc", feature = "transport-moq")
))]
pub(crate) async fn activate_native_iroh_packet_carrier(
&self,
kind: IrohPathKind,
remote_endpoint_id: iroh::EndpointId,
) -> anyhow::Result<crate::packet_carrier_transport::PacketCarrierSession> {
let carrier = match kind {
IrohPathKind::WebRtc => crate::iroh_carrier_kind::IrohCarrierKind::WebRtc,
IrohPathKind::Moq => crate::iroh_carrier_kind::IrohCarrierKind::Moq,
_ => anyhow::bail!("{kind:?} is not an Iroh packet carrier"),
};
self.activate_iroh_packet_carrier(carrier, remote_endpoint_id)
.await
}
pub async fn current_node_id(&self) -> Option<String> {
let endpoint_node_id = self
.iroh_endpoint
.read()
.await
.as_ref()
.map(|endpoint| endpoint.id().to_string());
let Some(endpoint_node_id) = endpoint_node_id else {
return self.node_id.read().await.clone();
};
let mut node_guard = self.node_id.write().await;
if node_guard.as_deref() != Some(endpoint_node_id.as_str()) {
eprintln!(
"[pluto-rtc][native] repaired local node projection from live endpoint old_node_id={} endpoint_node_id={}",
node_guard.as_deref().unwrap_or("<none>"),
endpoint_node_id
);
*node_guard = Some(endpoint_node_id.clone());
}
Some(endpoint_node_id)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn adopt_endpoint(&self, endpoint: Endpoint) {
if let Err(error) = self.adopt_endpoint_with_router_mode(endpoint, true).await {
eprintln!("[pluto-rtc][native] failed to adopt iroh endpoint: {error}");
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn adopt_endpoint_with_router_mode(
&self,
endpoint: Endpoint,
spawn_internal_router: bool,
) -> anyhow::Result<String> {
let node_id = endpoint.id().to_string();
let _init_lock = self.iroh_init_guard.lock().await;
{
let endpoint_guard = self.iroh_endpoint.read().await;
if let Some(existing) = endpoint_guard.as_ref() {
let existing_node_id = existing.id().to_string();
drop(endpoint_guard);
if existing_node_id == node_id {
*self.node_id.write().await = Some(existing_node_id.clone());
return Ok(existing_node_id);
}
endpoint.close().await;
anyhow::bail!(
"native Iroh endpoint already initialized as {existing_node_id}; refusing replacement with {node_id}"
);
}
}
let mut node_guard = self.node_id.write().await;
*node_guard = Some(node_id.clone());
drop(node_guard);
let mut endpoint_guard = self.iroh_endpoint.write().await;
*endpoint_guard = Some(endpoint);
let endpoint = endpoint_guard.clone().ok_or_else(|| {
anyhow::anyhow!("adopted endpoint disappeared before runtime install")
})?;
drop(endpoint_guard);
let node = if spawn_internal_router {
IrohNativeNode::spawn_with_endpoint(endpoint).await?
} else {
IrohNativeNode::spawn_with_endpoint_no_router(endpoint).await?
};
let native_accept_events = node.accept_events();
let native_incoming_streams = node.incoming_streams_stream();
*self.iroh_node.write().await = Some(node);
if spawn_internal_router {
self.start_native_accept_bridge(native_accept_events);
} else {
self.start_native_external_outbound_close_bridge(native_accept_events);
}
self.start_native_incoming_stream_router(native_incoming_streams);
eprintln!(
"[pluto-rtc][native] adopted endpoint node_id={} internal_router={}",
node_id, spawn_internal_router
);
Ok(node_id)
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn register_native_custom_transport_kind(
&self,
transport_id: u64,
kind: IrohPathKind,
) -> anyhow::Result<()> {
if transport_id == 0 {
anyhow::bail!("custom transport id must be non-zero");
}
let route = match kind {
IrohPathKind::Ble => crate::route_policy::KnownRoute::Ble,
IrohPathKind::WebRtc => crate::route_policy::KnownRoute::WebRtc,
IrohPathKind::Moq => crate::route_policy::KnownRoute::Moq,
_ => anyhow::bail!("unsupported custom transport path kind: {kind:?}"),
};
if route.descriptor().family != crate::route_policy::RouteFamily::IrohPhysical {
anyhow::bail!("custom transport route must use the iroh-physical family");
}
self.native_custom_transport_kinds
.write()
.await
.insert(transport_id, kind);
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn add_upgrade_provider(
&self,
provider: Arc<dyn crate::client::UpgradeProvider>,
) -> anyhow::Result<()> {
let kind = provider.kind();
let transport_id = provider.transport_id();
self.register_native_custom_transport_kind(transport_id, kind)
.await?;
self.native_transport_upgrade_providers
.write()
.await
.insert(kind, provider);
if kind == IrohPathKind::Ble {
self.wake_native_ble_recovery("provider-registered").await;
}
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn native_transport_upgrade_provider(
&self,
kind: IrohPathKind,
) -> Option<Arc<dyn crate::client::UpgradeProvider>> {
self.native_transport_upgrade_providers
.read()
.await
.get(&kind)
.cloned()
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn native_transport_upgrade_available(&self, kind: IrohPathKind) -> bool {
self.native_transport_upgrade_providers
.read()
.await
.contains_key(&kind)
}
pub async fn node_addr(&self) -> anyhow::Result<iroh::EndpointAddr> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.node_addr().await
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
async fn relay_only_mode_enabled(&self) -> bool {
#[cfg(not(target_arch = "wasm32"))]
{
should_use_relay_only_mode() || self.transport_config.read().await.iroh_relay_only
}
#[cfg(target_arch = "wasm32")]
{
true
}
}
pub(crate) fn relay_only_endpoint_addr(
mut addr: iroh::EndpointAddr,
) -> anyhow::Result<iroh::EndpointAddr> {
addr.addrs
.retain(|transport_addr| transport_addr.is_relay());
if addr.addrs.is_empty() {
anyhow::bail!(
"relay-only endpoint ticket unavailable: no relay address is currently online"
);
}
Ok(addr)
}
#[cfg(not(target_arch = "wasm32"))]
async fn relay_only_endpoint_addr_with_retry(
&self,
mut addr: iroh::EndpointAddr,
) -> anyhow::Result<iroh::EndpointAddr> {
let started_at = std::time::Instant::now();
let mut attempt: u32 = 0;
loop {
match Self::relay_only_endpoint_addr(addr.clone()) {
Ok(relay_addr) => return Ok(relay_addr),
Err(error) => {
if started_at.elapsed() >= RELAY_ONLY_ENDPOINT_TICKET_RETRY_WINDOW {
return Err(error);
}
}
}
attempt = attempt.saturating_add(1);
let multiplier = 1u32
.checked_shl(attempt.saturating_sub(1))
.unwrap_or(u32::MAX);
let delay = RELAY_ONLY_ENDPOINT_TICKET_MIN_RETRY_DELAY
.saturating_mul(multiplier)
.min(RELAY_ONLY_ENDPOINT_TICKET_MAX_RETRY_DELAY);
tokio::time::sleep(delay).await;
let endpoint = { self.iroh_endpoint.read().await.clone() };
if let Some(endpoint) = endpoint {
addr = endpoint.watch_addr().get();
} else {
addr = self.node_addr().await?;
}
}
}
pub async fn endpoint_ticket(&self) -> anyhow::Result<String> {
let mut addr = self.node_addr().await?;
if self.relay_only_mode_enabled().await {
#[cfg(not(target_arch = "wasm32"))]
{
addr = self.relay_only_endpoint_addr_with_retry(addr).await?;
}
#[cfg(target_arch = "wasm32")]
{
addr = Self::relay_only_endpoint_addr(addr)?;
}
}
Ok(EndpointTicket::new(addr).to_string())
}
#[cfg(not(target_arch = "wasm32"))]
async fn persistent_endpoint_ticket_with_token(
&self,
scope: &str,
max_connections: u32,
) -> anyhow::Result<String> {
let latest_iroh_ticket = self.endpoint_ticket().await?;
let normalized_scope = scope.trim();
let cached_result = {
let mut cache = match self.managed_scope_tickets.write() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
if let Some(existing) = cache.get_mut(normalized_scope) {
let mut persist_after_release: Option<(String, u32)> = None;
let compound_ticket = if existing.max_connections != max_connections {
self.session_token_registry.revoke(existing.token.as_str());
let token = crate::session_token::generate_token();
let grant_scope =
crate::session_token::GrantScope::from(normalized_scope.to_string());
self.session_token_registry.register(
token.clone(),
grant_scope.clone(),
max_connections,
);
existing.scope = grant_scope;
existing.token = token;
existing.max_connections = max_connections;
existing.iroh_ticket = latest_iroh_ticket.clone();
existing.compound_ticket = crate::session_token::build_ticket(
&latest_iroh_ticket,
&existing.token,
&existing.scope,
existing.max_connections,
);
persist_after_release =
Some((existing.token.clone(), existing.max_connections));
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(existing.compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][managed-cache-reset] scope={} token_fp={} iroh_fp={} max_connections={} reason=max-connections-changed",
scope_name.unwrap_or_else(|| existing.scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint,
max_connections
);
existing.compound_ticket.clone()
} else if existing.iroh_ticket != latest_iroh_ticket {
existing.iroh_ticket = latest_iroh_ticket.clone();
existing.compound_ticket = crate::session_token::build_ticket(
&latest_iroh_ticket,
&existing.token,
&existing.scope,
existing.max_connections,
);
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(existing.compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][managed-cache-refresh] scope={} token_fp={} iroh_fp={} endpoint_changed=true",
scope_name.unwrap_or_else(|| existing.scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint
);
existing.compound_ticket.clone()
} else {
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(existing.compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][managed-cache-reuse] scope={} token_fp={} iroh_fp={} max_connections={}",
scope_name.unwrap_or_else(|| existing.scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint,
existing.max_connections
);
existing.compound_ticket.clone()
};
Some((compound_ticket, persist_after_release))
} else {
None
}
};
if let Some((compound_ticket, persist_after_release)) = cached_result {
if let Some((token, persisted_max_connections)) = persist_after_release {
self.persist_managed_scope_grant(
normalized_scope,
token.as_str(),
persisted_max_connections,
)
.await?;
}
return Ok(compound_ticket);
}
if let Some(persisted) = self
.load_persisted_managed_scope_grant(normalized_scope)
.await?
{
if persisted.max_connections == max_connections {
let grant_scope = crate::session_token::GrantScope::from(persisted.scope.clone());
self.session_token_registry.register(
persisted.token.clone(),
grant_scope.clone(),
persisted.max_connections,
);
let compound_ticket = crate::session_token::build_ticket(
&latest_iroh_ticket,
&persisted.token,
&grant_scope,
persisted.max_connections,
);
let mut cache = match self.managed_scope_tickets.write() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
cache.insert(
normalized_scope.to_string(),
CachedManagedScopeTicket {
scope: grant_scope.clone(),
token: persisted.token,
max_connections: persisted.max_connections,
compound_ticket: compound_ticket.clone(),
iroh_ticket: latest_iroh_ticket,
},
);
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][managed-persist-load] scope={} token_fp={} iroh_fp={} max_connections={}",
scope_name.unwrap_or_else(|| grant_scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint,
max_connections
);
return Ok(compound_ticket);
}
self.session_token_registry.revoke(persisted.token.as_str());
self.delete_persisted_managed_scope_grant(normalized_scope)
.await?;
eprintln!(
"[PlutoRTC][ticket][managed-persist-reset] scope={} reason=max-connections-changed",
normalized_scope
);
}
let token = crate::session_token::generate_token();
let grant_scope = crate::session_token::GrantScope::from(normalized_scope.to_string());
self.session_token_registry
.register(token.clone(), grant_scope.clone(), max_connections);
let compound_ticket = crate::session_token::build_ticket(
&latest_iroh_ticket,
&token,
&grant_scope,
max_connections,
);
{
let mut cache = match self.managed_scope_tickets.write() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
cache.insert(
normalized_scope.to_string(),
CachedManagedScopeTicket {
scope: grant_scope.clone(),
token: token.clone(),
max_connections,
compound_ticket: compound_ticket.clone(),
iroh_ticket: latest_iroh_ticket,
},
);
}
self.persist_managed_scope_grant(normalized_scope, token.as_str(), max_connections)
.await?;
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][managed-cache-create] scope={} token_fp={} iroh_fp={} max_connections={}",
scope_name.unwrap_or_else(|| grant_scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint,
max_connections
);
Ok(compound_ticket)
}
pub async fn endpoint_ticket_with_token(
&self,
scope: &str,
max_connections: u32,
) -> anyhow::Result<String> {
#[cfg(not(target_arch = "wasm32"))]
if scope.trim() == PERSISTENT_MANAGED_ADMISSION_SCOPE {
return self
.persistent_endpoint_ticket_with_token(scope, max_connections)
.await;
}
let iroh_ticket = self.endpoint_ticket().await?;
let token = crate::session_token::generate_token();
let grant_scope = crate::session_token::GrantScope::from(scope);
let expires_at_ms = if scope.trim() == "user-device" {
None
} else {
Some(
crate::session_token::now_unix_ms()
.saturating_add(crate::session_token::DEFAULT_RESTRICTED_SESSION_TOKEN_TTL_MS),
)
};
self.session_token_registry.register_with_expiry_ms(
token.clone(),
grant_scope.clone(),
max_connections,
expires_at_ms,
);
let compound_ticket = crate::session_token::expiring_ticket(
&iroh_ticket,
&token,
&grant_scope,
max_connections,
expires_at_ms,
);
let (iroh_fingerprint, scope_name, token_fingerprint) =
summarize_compound_ticket_for_logs(compound_ticket.as_str());
eprintln!(
"[PlutoRTC][ticket][mint] scope={} token_fp={} iroh_fp={} max_connections={}",
scope_name.unwrap_or_else(|| grant_scope.to_string()),
token_fingerprint.unwrap_or_else(|| "unknown".to_string()),
iroh_fingerprint,
max_connections
);
Ok(compound_ticket)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn export_endpoint_handle(&self) -> anyhow::Result<EndpointHandle> {
let node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let node_addr = self.node_addr().await?;
Ok(EndpointHandle {
node_id,
node_addr: serde_json::to_string(&node_addr)
.map_err(|e| anyhow::anyhow!("failed to serialize node addr: {}", e))?,
})
}
async fn endpoint_dial_attempt_fences(
&self,
endpoint_id: &iroh::EndpointId,
) -> Vec<(String, u64, Option<u64>)> {
let eid = endpoint_id.to_string();
let mut records = self.connection_manager.get_by_endpoint_id(&eid).await;
if records.is_empty() {
records = self.connection_manager.get_by_node_id(&eid).await;
}
records
.into_iter()
.filter(|record| {
matches!(
record.state,
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting
)
})
.map(|record| {
(
record.connection_id,
record.transport_generation,
record.transport_stable_id,
)
})
.collect()
}
async fn mark_endpoint_dial_timed_out(
&self,
attempt_fences: &[(String, u64, Option<u64>)],
label: &'static str,
) {
let reason = dial_timeout_reason(label);
for (connection_id, transport_generation, transport_stable_id) in attempt_fences {
let _ = self
.connection_manager
.set_failed_if_dial_attempt_current(
connection_id,
*transport_generation,
*transport_stable_id,
Some(reason.clone()),
)
.await;
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn connect(&self, endpoint_id: iroh::EndpointId) -> anyhow::Result<BiStream> {
self.ensure_connected(endpoint_id).await?;
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
let (send, recv) = node.open_bi(endpoint_id).await?;
let connection_id = self
.connection_manager
.get_by_endpoint_id(&endpoint_id.to_string())
.await
.into_iter()
.next()
.map(|record| record.connection_id);
let key = self.application_crypto_key_for_connection(connection_id.as_deref());
let (send, recv) = match connection_id.as_deref() {
Some(id) => self.wrap_application_streams(id, key, send, recv, &[])?,
None => crate::application_crypto_streams::wrap_peer_streams(None, send, recv)?,
};
Ok(BiStream {
send,
recv,
id: endpoint_id.to_string(),
})
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn ensure_connected(&self, endpoint_id: iroh::EndpointId) -> anyhow::Result<()> {
self.ensure_connected_with_timeout(endpoint_id, std::time::Duration::from_secs(4))
.await
}
pub(crate) async fn remember_endpoint_addr(&self, endpoint_addr: &iroh::EndpointAddr) {
self.known_endpoint_addrs
.write()
.await
.insert(endpoint_addr.id.to_string(), endpoint_addr.clone());
self.known_device_endpoint_revision
.send_modify(|revision| *revision = revision.wrapping_add(1));
}
pub(crate) async fn cached_endpoint_addr(
&self,
endpoint_id: iroh::EndpointId,
) -> Option<iroh::EndpointAddr> {
self.known_endpoint_addrs
.read()
.await
.get(&endpoint_id.to_string())
.cloned()
}
pub(crate) async fn finalize_transport_dial_record(
&self,
endpoint_id: iroh::EndpointId,
transport_source: Option<String>,
) -> anyhow::Result<String> {
let remote_node_id = endpoint_id.to_string();
let local_node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
if !self.is_connection_transport_alive(endpoint_id).await {
return Ok(connection_id);
}
let device_id_hint = self
.known_remote_device_id_for_incoming_transport(&connection_id, &remote_node_id)
.await;
const MAX_COMMIT_FENCE_PASSES: usize = 3;
let mut source = transport_source;
for pass in 0..MAX_COMMIT_FENCE_PASSES {
let Some(transport_stable_id) = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection))
else {
return Ok(connection_id);
};
self.commit_current_native_transport_record(
&connection_id,
&remote_node_id,
device_id_hint.clone(),
Some(transport_stable_id),
source.clone(),
)
.await;
let committed_transport_is_current = self
.get_connection(endpoint_id)
.await
.is_some_and(|connection| {
crate::transport_generation::for_connection(&connection) == transport_stable_id
});
if committed_transport_is_current {
return Ok(connection_id);
}
eprintln!(
"[PlutoRTC] transport commit fence retry connection_id={} remote_node_id={} stale_stable_id={} pass={}",
connection_id,
remote_node_id,
transport_stable_id,
pass + 1,
);
source = Some("physical-arbitration-reconcile".to_string());
}
eprintln!(
"[PlutoRTC] transport commit fence exhausted connection_id={} remote_node_id={}; next registry event will reconcile",
connection_id, remote_node_id,
);
Ok(connection_id)
}
pub(crate) async fn commit_current_native_transport_record(
&self,
connection_id: &str,
remote_node_id: &str,
device_id_hint: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
) -> bool {
self.commit_current_native_transport_record_inner(
connection_id,
remote_node_id,
device_id_hint,
None,
transport_stable_id,
transport_source,
None,
crate::connection_manager::TerminalReopenAuthority::None,
|| true,
|_| {},
)
.await
}
#[cfg(target_arch = "wasm32")]
async fn commit_current_desired_user_device_transport_record(
&self,
connection_id: &str,
remote_node_id: &str,
device_id: &str,
transport_stable_id: u64,
transport_source: Option<String>,
) -> bool {
let expected_device_id = device_id.to_string();
self.commit_current_native_transport_record_inner(
connection_id,
remote_node_id,
Some(expected_device_id.clone()),
None,
Some(transport_stable_id),
transport_source,
None,
crate::connection_manager::TerminalReopenAuthority::CurrentDesiredUserDevice,
|| {
self.current_desired_user_device_for_node_now(remote_node_id)
.as_deref()
== Some(expected_device_id.as_str())
},
|_| {},
)
.await
}
pub(crate) async fn commit_current_admitted_native_transport_record(
&self,
connection_id: &str,
remote_node_id: &str,
device_id_hint: Option<String>,
authoritative_device_id: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
authenticated_peer_device_id: Option<&str>,
commit_application_security: impl FnOnce() -> bool,
normalize_scopes: impl FnOnce(&mut std::collections::HashSet<String>),
) -> bool {
self.commit_current_native_transport_record_inner(
connection_id,
remote_node_id,
device_id_hint,
authoritative_device_id,
transport_stable_id,
transport_source,
authenticated_peer_device_id,
crate::connection_manager::TerminalReopenAuthority::FreshSessionCapability,
commit_application_security,
normalize_scopes,
)
.await
}
async fn commit_current_native_transport_record_inner(
&self,
connection_id: &str,
remote_node_id: &str,
device_id_hint: Option<String>,
authoritative_device_id: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
authenticated_peer_device_id: Option<&str>,
terminal_reopen_authority: crate::connection_manager::TerminalReopenAuthority,
commit_application_security: impl FnOnce() -> bool,
normalize_scopes: impl FnOnce(&mut std::collections::HashSet<String>),
) -> bool {
let Some(transport_stable_id) = transport_stable_id else {
return false;
};
let Ok(remote_endpoint_id) = remote_node_id.parse::<iroh::EndpointId>() else {
eprintln!(
"[OpenRTC][transport] rejected physical commit with invalid endpoint connection_id={} remote_node_id={}",
connection_id, remote_node_id,
);
return false;
};
let node = self.iroh_node.read().await.as_ref().cloned();
let Some(node) = node else {
return false;
};
let mut commit_application_security = Some(commit_application_security);
let commit_result = node
.with_current_connection_generation(
remote_endpoint_id,
transport_stable_id,
|| async {
if self
.connection_manager
.get_by_connection_id(connection_id)
.await
.is_none()
{
if self
.auto_connect_exclusion_for_peer(
device_id_hint.as_deref(),
Some(remote_node_id),
)
.is_some()
{
let _ = self
.connection_manager
.materialize_manual_disconnect_tombstone(
connection_id.to_string(),
Some(remote_node_id.to_string()),
device_id_hint.clone(),
Some(remote_node_id.to_string()),
)
.await;
if authenticated_peer_device_id.is_none() {
return None;
}
} else {
let _ = self
.connection_manager
.upsert_pending_for_physical_observation(
connection_id.to_string(),
Some(remote_node_id.to_string()),
device_id_hint.clone(),
Some(remote_node_id.to_string()),
)
.await;
}
}
#[cfg(not(target_arch = "wasm32"))]
self.retire_stale_native_main_route(connection_id, transport_stable_id)
.await;
let mut before_publish = |previous_transport_stable_id| {
let Some(commit) = commit_application_security.take() else {
return false;
};
if !commit() {
return false;
}
if transport_replacement_requires_application_crypto_reconfirmation(
previous_transport_stable_id,
transport_stable_id,
) {
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[OpenRTC][transport] physical replacement requires fresh application-crypto confirmation connection_id={} previous_stable_id={:?} current_stable_id={}",
connection_id, previous_transport_stable_id, transport_stable_id,
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[OpenRTC][transport] physical replacement requires fresh application-crypto confirmation connection_id={} previous_stable_id={:?} current_stable_id={}",
connection_id, previous_transport_stable_id, transport_stable_id,
)));
}
self.bind_connection_application_crypto_transport(
connection_id,
transport_stable_id,
);
true
};
self.connection_manager
.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
connection_id,
Some(remote_node_id.to_string()),
authoritative_device_id.clone(),
authoritative_device_id,
Some(transport_stable_id),
transport_source,
|requires_authenticated_peer_reconnect, previous_transport_stable_id| {
if !requires_authenticated_peer_reconnect {
return before_publish(previous_transport_stable_id);
}
let Some(device_id) = authenticated_peer_device_id else {
return false;
};
let resumed = self.commit_peer_requested_reconnect_if_authenticated(
device_id,
|| before_publish(previous_transport_stable_id),
);
if resumed {
eprintln!(
"[PlutoRTC][session-admission][peer-reconnect] atomically committed security and cleared peer-requested auto-connect suppression connection_id={} device_id={}",
connection_id, device_id,
);
}
resumed
},
normalize_scopes,
terminal_reopen_authority,
)
.await
},
)
.await;
let committed = match commit_result {
Some(Some(record)) => record,
None => {
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[OpenRTC][transport] rejected physical commit connection_id={} remote_node_id={} stable_id={} reason=physical-generation-not-current",
connection_id, remote_node_id, transport_stable_id,
);
return false;
}
Some(None) => {
#[cfg(not(target_arch = "wasm32"))]
{
let record = self
.connection_manager
.get_by_connection_id(connection_id)
.await;
eprintln!(
"[OpenRTC][transport] rejected physical commit connection_id={} remote_node_id={} stable_id={} reason=logical-or-security-denied manager_state={:?} manager_stable_id={:?} manager_status_reason={:?} manager_last_disconnect_reason={:?}",
connection_id,
remote_node_id,
transport_stable_id,
record.as_ref().map(|record| &record.state),
record.as_ref().and_then(|record| record.transport_stable_id),
record.as_ref().and_then(|record| record.status_reason.as_deref()),
record
.as_ref()
.and_then(|record| record.last_disconnect_reason.as_deref()),
);
}
return false;
}
};
if committed.transport_stable_id != Some(transport_stable_id) {
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[OpenRTC][transport] rejected physical commit connection_id={} remote_node_id={} stable_id={} reason=manager-returned-different-generation manager_stable_id={:?}",
connection_id, remote_node_id, transport_stable_id, committed.transport_stable_id,
);
return false;
}
#[cfg(target_arch = "wasm32")]
{
let active_transport = self.iroh_path_kind(remote_node_id).await.transport_label();
let _ = self
.report_transport_status_for_current_generation(
connection_id,
active_transport,
None,
)
.await;
}
true
}
#[cfg(feature = "iroh-carrier-core")]
pub(crate) async fn begin_crypto_bound_transport_replacement_commit(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
expected_route_generation: u64,
replacement_transport_stable_id: u64,
transport_source: Option<String>,
) -> Option<crate::connection_manager::TransportReplacementCommit> {
self.connection_manager
.begin_transport_replacement_commit_if_generation_current(
connection_id,
expected_transport_stable_id,
expected_transport_generation,
expected_route_generation,
replacement_transport_stable_id,
transport_source,
|_| {
self.bind_connection_application_crypto_transport(
connection_id,
replacement_transport_stable_id,
);
},
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn promote_known_native_user_device_connection(
&self,
connection_id: &str,
remote_node_id: &str,
) -> Option<String> {
let device_id = self
.known_remote_device_id_for_incoming_transport(connection_id, remote_node_id)
.await?;
if !self
.mark_trusted_user_device_connection_admitted(connection_id, &device_id)
.await
{
return None;
}
let transport_stable_id = self
.connection_manager
.get_by_connection_id(connection_id)
.await
.and_then(|record| record.transport_stable_id);
if self.trusted_user_device_application_crypto_is_required()
&& !self.connection_application_crypto_is_confirmed(connection_id, transport_stable_id)
{
self.set_connection_application_crypto_required(connection_id);
match self.get_or_create_connection_key_agreement(connection_id) {
Ok(key_agreement) => {
if let Err(error) = self
.send_typescript_capability_update(
connection_id,
"trusted-user-device-application-crypto",
Some(key_agreement.public_key_bytes()),
)
.await
{
eprintln!(
"[OpenRTC][capability] trusted user-device key agreement seed failed connection_id={} error={}",
connection_id, error,
);
}
}
Err(error) => {
eprintln!(
"[OpenRTC][capability] trusted user-device key agreement initialization failed connection_id={} error={:?}",
connection_id, error,
);
}
}
}
let _ = self
.confirm_managed_connection_readiness(connection_id)
.await;
Some(device_id)
}
pub(crate) async fn ensure_connection_manager_record_before_peer_stream(
&self,
endpoint_id: &iroh::EndpointId,
) -> anyhow::Result<()> {
use crate::connection_manager::ConnectionState;
if !self.is_connection_transport_alive(*endpoint_id).await {
return Ok(());
}
let remote_node_id = endpoint_id.to_string();
let Some(local_node_id) = self.current_node_id().await else {
return Ok(());
};
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
let physical_transport_stable_id = self
.get_connection(*endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let needs_commit = match self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
None => true,
Some(record) => {
!matches!(record.state, ConnectionState::Connected)
|| record.transport_stable_id != physical_transport_stable_id
}
};
if needs_commit {
self.finalize_transport_dial_record(*endpoint_id, Some("peer-stream-open".to_string()))
.await?;
}
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
#[doc(hidden)]
pub async fn ensure_connected_with_timeout(
&self,
endpoint_id: iroh::EndpointId,
timeout: std::time::Duration,
) -> anyhow::Result<()> {
use futures::StreamExt;
let already_active = {
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
node.is_connected(endpoint_id).await
};
if already_active {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected-active".to_string()),
)
.await?;
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
self.spawn_outgoing_path_watcher_if_available(endpoint_id, node)
.await;
}
return Ok(());
}
if let Some(endpoint_addr) = self.cached_endpoint_addr(endpoint_id).await {
eprintln!(
"[pluto-rtc][iroh] redial using cached endpoint addr endpoint_id={}",
endpoint_id
);
return self.ensure_connected_addr(endpoint_id, endpoint_addr).await;
}
let dial_attempt_fences = self.endpoint_dial_attempt_fences(&endpoint_id).await;
let mut events = {
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
node.connect(endpoint_id)
};
let timeout = tokio::time::sleep(timeout);
tokio::pin!(timeout);
loop {
tokio::select! {
_ = &mut timeout => {
self.mark_endpoint_dial_timed_out(&dial_attempt_fences, "ensure_connected")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next() => {
match event {
Some(ConnectEvent::Connected) => {
let connection_id = self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected".to_string()),
)
.await?;
self.restore_scoped_route_admission(endpoint_id, &connection_id)
.await?;
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
self.spawn_outgoing_path_watcher_if_available(endpoint_id, node)
.await;
return Ok(());
}
Some(ConnectEvent::Closed { error }) => {
if let Some(error) = error {
return Err(anyhow::anyhow!("Failed to connect to {}: {}", endpoint_id, error));
}
return Err(anyhow::anyhow!("Connection to {} closed before it was established", endpoint_id));
}
None => {
return Err(anyhow::anyhow!("Connection stream ended before connecting to {}", endpoint_id));
}
}
}
}
}
}
#[cfg(target_arch = "wasm32")]
pub async fn ensure_connected(&self, endpoint_id: iroh::EndpointId) -> anyhow::Result<()> {
use futures::{FutureExt, StreamExt};
use std::time::Duration;
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
if node.is_connected(endpoint_id).await {
drop(node_guard);
let connection_id = self
.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected-active".to_string()),
)
.await?;
self.emit_current_wasm_connection_state(&connection_id)
.await;
return Ok(());
}
if let Some(endpoint_addr) = self.cached_endpoint_addr(endpoint_id).await {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][iroh] redial using cached endpoint addr endpoint_id={}",
endpoint_id
)));
return self.ensure_connected_addr(endpoint_id, endpoint_addr).await;
}
let dial_attempt_fences = self.endpoint_dial_attempt_fences(&endpoint_id).await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][core_impl][ensure_connected] starting dial endpoint_id={}",
endpoint_id
)));
let mut events = node.connect(endpoint_id);
let timeout = gloo_timers::future::sleep(Duration::from_secs(15)).fuse();
futures::pin_mut!(timeout);
loop {
futures::select! {
_ = timeout => {
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][core_impl][ensure_connected] TIMEOUT — dropping events receiver endpoint_id={}. The underlying connect task may still succeed but will have no one to send ConnectEvent to.",
endpoint_id
)));
self.mark_endpoint_dial_timed_out(&dial_attempt_fences, "ensure_connected")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next().fuse() => {
match event {
Some(ConnectEvent::Connected) => {
let connection_id = self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected".to_string()),
)
.await?;
self.restore_scoped_route_admission(endpoint_id, &connection_id)
.await?;
self.emit_current_wasm_connection_state(&connection_id).await;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| {
crate::transport_generation::for_connection(
&connection,
)
});
Arc::new(self.clone())
.start_wasm_connect_event_bridge(
endpoint_id,
transport_stable_id,
events,
);
return Ok(());
}
Some(ConnectEvent::Closed { error }) => {
if let Some(error) = error {
return Err(anyhow::anyhow!("Failed to connect to {}: {}", endpoint_id, error));
}
return Err(anyhow::anyhow!("Connection to {} closed before it was established", endpoint_id));
}
None => {
return Err(anyhow::anyhow!("Connection stream ended before connecting to {}", endpoint_id));
}
}
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn ensure_connected_addr(
&self,
endpoint_id: iroh::EndpointId,
endpoint_addr: iroh::EndpointAddr,
) -> anyhow::Result<()> {
use futures::StreamExt;
use std::time::Duration;
let already_active = {
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
node.is_connected(endpoint_id).await
};
if already_active {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr-active".to_string()),
)
.await?;
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
self.spawn_outgoing_path_watcher_if_available(endpoint_id, node)
.await;
}
return Ok(());
}
self.remember_endpoint_addr(&endpoint_addr).await;
let dial_attempt_fences = self.endpoint_dial_attempt_fences(&endpoint_id).await;
let mut events = {
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
node.connect_addr(endpoint_id, endpoint_addr)
};
let timeout = tokio::time::sleep(Duration::from_secs(4));
tokio::pin!(timeout);
loop {
tokio::select! {
_ = &mut timeout => {
self.mark_endpoint_dial_timed_out(&dial_attempt_fences, "ensure_connected_addr")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next() => {
match event {
Some(ConnectEvent::Connected) => {
let connection_id = self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr".to_string()),
)
.await?;
self.restore_scoped_route_admission(endpoint_id, &connection_id)
.await?;
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
self.spawn_outgoing_path_watcher_if_available(endpoint_id, node)
.await;
return Ok(());
}
Some(ConnectEvent::Closed { error }) => {
if let Some(error) = error {
return Err(anyhow::anyhow!("Failed to connect to {}: {}", endpoint_id, error));
}
return Err(anyhow::anyhow!("Connection to {} closed before it was established", endpoint_id));
}
None => {
return Err(anyhow::anyhow!("Connection stream ended before connecting to {}", endpoint_id));
}
}
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn spawn_outgoing_path_watcher_if_available(
&self,
endpoint_id: iroh::EndpointId,
node: &IrohNativeNode,
) {
let Some(local_node_id) = self.current_node_id().await else {
return;
};
let remote_node_id = endpoint_id.to_string();
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
if let Some(connection) = node.get_connection(endpoint_id).await {
self.spawn_iroh_path_watcher(&connection_id, &remote_node_id, &connection, false);
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn reconcile_native_accept_close(
&self,
endpoint_id: iroh::EndpointId,
transport_stable_id: u64,
error: Option<String>,
) {
let endpoint_id_str = endpoint_id.to_string();
let recoverable_close = !is_terminal_close_reason(error.as_deref());
let mut records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
if let Some(local_node_id) = self.current_node_id().await {
let deterministic_connection_id =
Self::deterministic_connection_id(&local_node_id, &endpoint_id_str);
if !records
.iter()
.any(|record| record.connection_id == deterministic_connection_id)
{
if let Some(record) = self
.connection_manager
.get_by_connection_id(&deterministic_connection_id)
.await
{
eprintln!(
"[pluto-rtc][native-accept] endpoint index fallback connection_id={} endpoint_id={} transport_stable_id={}",
deterministic_connection_id,
endpoint_id,
transport_stable_id,
);
records.push(record);
}
}
}
let mut replacement_observed = false;
for record in records {
if self
.observe_native_transport_generation_closed(
&record.connection_id,
transport_stable_id,
error.as_deref(),
"native-accept-closed",
)
.await
{
replacement_observed = true;
eprintln!(
"[pluto-rtc][native-accept] current transport closed connection_id={} endpoint_id={} transport_stable_id={} error={:?}",
record.connection_id,
endpoint_id,
transport_stable_id,
error,
);
self.emit_current_native_connection_state(&record.connection_id)
.await;
}
}
if replacement_observed && recoverable_close {
let _ = self
.schedule_native_external_auto_connect_recovery(&endpoint_id_str)
.await;
}
}
#[cfg(not(target_arch = "wasm32"))]
fn start_native_external_outbound_close_bridge(
&self,
mut events: futures::stream::BoxStream<'static, AcceptEvent>,
) {
let client = self.clone();
tokio::spawn(async move {
use futures::StreamExt as _;
while let Some(event) = events.next().await {
let AcceptEvent::Closed {
endpoint_id,
transport_stable_id,
error,
was_outbound: true,
..
} = event
else {
continue;
};
client
.reconcile_native_accept_close(endpoint_id, transport_stable_id, error)
.await;
}
});
}
#[cfg(not(target_arch = "wasm32"))]
fn start_native_accept_bridge(
&self,
mut events: futures::stream::BoxStream<'static, AcceptEvent>,
) {
let client = self.clone();
tokio::spawn(async move {
use futures::StreamExt as _;
while let Some(event) = events.next().await {
let (endpoint_id, transport_stable_id, replaced_transport_stable_id) = match event {
AcceptEvent::Accepted {
endpoint_id,
transport_stable_id,
replaced_transport_stable_id,
} => (
endpoint_id,
transport_stable_id,
replaced_transport_stable_id,
),
AcceptEvent::Closed {
endpoint_id,
transport_stable_id,
error,
..
} => {
client
.reconcile_native_accept_close(endpoint_id, transport_stable_id, error)
.await;
continue;
}
};
let current_stable_id = client
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if !accepted_transport_event_is_current(transport_stable_id, current_stable_id) {
eprintln!(
"[pluto-rtc][native-accept] ignored stale accepted transport endpoint_id={} accepted_stable_id={} current_stable_id={:?}",
endpoint_id, transport_stable_id, current_stable_id
);
continue;
}
if let Some(withdrawn_device_id) = client
.withdrawn_desired_device_for_node(&endpoint_id.to_string())
.await
{
let _ = client
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_DESIRED_PEER_WITHDRAWN,
)
.await;
eprintln!(
"[pluto-rtc][native-accept] rejected withdrawn desired peer endpoint_id={} device_id={} transport_stable_id={}",
endpoint_id, withdrawn_device_id, transport_stable_id,
);
continue;
}
let Some(local_node_id) = client.current_node_id().await else {
eprintln!(
"[pluto-rtc][native-accept] missing local node while classifying accepted transport endpoint_id={}",
endpoint_id,
);
continue;
};
let connection_id =
Client::deterministic_connection_id(&local_node_id, &endpoint_id.to_string());
let projected_transport_generation = client
.connection_manager
.peer_snapshot(&connection_id)
.await
.map(|snapshot| snapshot.active_transport_generation)
.unwrap_or_default();
let logical_session_was_admitted = matches!(
client.session_admission(&connection_id),
crate::session_token::SessionAdmission::Accepted { .. }
);
let accepted_replacement = accepted_transport_replaces_managed_generation(
replaced_transport_stable_id,
transport_stable_id,
projected_transport_generation,
logical_session_was_admitted,
);
let connection_id = match client
.finalize_transport_dial_record(endpoint_id, Some("native-accept".to_string()))
.await
{
Ok(connection_id) => connection_id,
Err(error) => {
eprintln!(
"[pluto-rtc][native-accept] failed projecting accepted transport endpoint_id={} error={}",
endpoint_id, error
);
continue;
}
};
let remote_node_id = endpoint_id.to_string();
if let Some(device_id) = client
.promote_known_native_user_device_connection(&connection_id, &remote_node_id)
.await
{
eprintln!(
"[pluto-rtc][native-accept] promoted trusted user-device connection_id={} remote_node_id={} device_id={}",
connection_id, remote_node_id, device_id
);
}
if accepted_replacement {
eprintln!(
"[pluto-rtc][native-accept][replacement-route] connection_id={} remote_node_id={} prior_stable_id={:?} accepted_stable_id={} projected_generation={} logical_session_was_admitted={}",
connection_id,
remote_node_id,
replaced_transport_stable_id,
transport_stable_id,
projected_transport_generation,
logical_session_was_admitted,
);
let _ = client
.wake_native_external_auto_connect_for_replaced_route(&remote_node_id)
.await;
}
client
.migrate_remote_admission_proof_for_pending_ble_upgrade(
&connection_id,
endpoint_id,
)
.await;
let node_guard = client.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
continue;
};
client
.spawn_outgoing_path_watcher_if_available(endpoint_id, node)
.await;
}
});
}
#[cfg(not(target_arch = "wasm32"))]
fn start_native_incoming_stream_router(
&self,
incoming: async_channel::Receiver<IncomingStream>,
) {
const MAX_IN_FLIGHT_NATIVE_STREAM_ROUTES: usize = 64;
let route_permits = Arc::new(tokio::sync::Semaphore::new(
MAX_IN_FLIGHT_NATIVE_STREAM_ROUTES,
));
let client = self.clone();
let application_streams = self.native_application_streams.clone();
tokio::spawn(async move {
while let Ok(incoming_stream) = incoming.recv().await {
if application_streams.is_closed() {
break;
}
let endpoint_id = incoming_stream.endpoint_id;
let transport_stable_id = incoming_stream.transport_stable_id;
match incoming_stream.stream {
crate::native_node::IncomingStreamType::Bi(mut send, recv) => {
let Ok(permit) = route_permits.clone().try_acquire_owned() else {
let _ = send.finish();
drop(recv);
eprintln!(
"[pluto-rtc][native-stream-router] dropped bi-stream because native admission ingress is saturated endpoint_id={} limit={}",
endpoint_id, MAX_IN_FLIGHT_NATIVE_STREAM_ROUTES,
);
continue;
};
let client = client.clone();
let application_streams = application_streams.clone();
tokio::spawn(async move {
let _permit = permit;
let disposition = client
.route_incoming_bi_stream_for_admission(
endpoint_id,
transport_stable_id,
send,
recv,
)
.await;
let disposition = match disposition {
Ok(crate::client::BiStreamDisposition::Forward {
send,
recv,
recv_prefix,
}) => client
.route_admitted_peer_message(
endpoint_id,
transport_stable_id,
send,
recv,
recv_prefix,
)
.await
.map_err(|error| format!("peer message intake: {error:#}")),
other => other,
};
match disposition {
Ok(crate::client::BiStreamDisposition::Consumed) => {}
Ok(crate::client::BiStreamDisposition::Forward {
send,
recv,
recv_prefix,
}) => {
if application_streams
.send(IncomingStream {
endpoint_id,
transport_stable_id,
recv_prefix,
stream: crate::native_node::IncomingStreamType::Bi(
send, recv,
),
})
.await
.is_err()
{
eprintln!(
"[pluto-rtc][native-stream-router] application stream receiver closed endpoint_id={}",
endpoint_id,
);
}
}
Err(error) => {
eprintln!(
"[pluto-rtc][native-stream-router] denied pre-admission bi-stream endpoint_id={} error={}",
endpoint_id, error
);
}
}
});
}
crate::native_node::IncomingStreamType::Uni(recv) => {
if client.native_peer_is_pending_admission(endpoint_id).await {
eprintln!(
"[pluto-rtc][native-stream-router] denied pending uni-stream endpoint_id={}",
endpoint_id
);
continue;
}
if application_streams
.send(IncomingStream {
endpoint_id,
transport_stable_id,
recv_prefix: Vec::new(),
stream: crate::native_node::IncomingStreamType::Uni(recv),
})
.await
.is_err()
{
break;
}
}
}
}
});
}
#[cfg(target_arch = "wasm32")]
pub async fn ensure_connected_addr(
&self,
endpoint_id: iroh::EndpointId,
endpoint_addr: iroh::EndpointAddr,
) -> anyhow::Result<()> {
use futures::{FutureExt, StreamExt};
use std::time::Duration;
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
if node.is_connected(endpoint_id).await {
drop(node_guard);
let connection_id = self
.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr-active".to_string()),
)
.await?;
self.emit_current_wasm_connection_state(&connection_id)
.await;
return Ok(());
}
self.remember_endpoint_addr(&endpoint_addr).await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][core_impl][ensure_connected_addr] starting dial endpoint_id={}",
endpoint_id
)));
let dial_attempt_fences = self.endpoint_dial_attempt_fences(&endpoint_id).await;
let mut events = node.connect_addr(endpoint_id, endpoint_addr);
let timeout = gloo_timers::future::sleep(Duration::from_secs(15)).fuse();
futures::pin_mut!(timeout);
loop {
futures::select! {
_ = timeout => {
web_sys::console::error_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][core_impl][ensure_connected_addr] TIMEOUT — dropping events receiver endpoint_id={}. The underlying connect_addr task may still succeed but will have no one to send ConnectEvent to.",
endpoint_id
)));
self.mark_endpoint_dial_timed_out(&dial_attempt_fences, "ensure_connected_addr")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next().fuse() => {
match event {
Some(ConnectEvent::Connected) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][core_impl][ensure_connected_addr] Connected event received endpoint_id={}",
endpoint_id
)));
let connection_id = self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr".to_string()),
)
.await?;
self.restore_scoped_route_admission(endpoint_id, &connection_id)
.await?;
self.emit_current_wasm_connection_state(&connection_id).await;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| {
crate::transport_generation::for_connection(
&connection,
)
});
Arc::new(self.clone())
.start_wasm_connect_event_bridge(
endpoint_id,
transport_stable_id,
events,
);
return Ok(());
}
Some(ConnectEvent::Closed { error }) => {
if let Some(error) = error {
return Err(anyhow::anyhow!("Failed to connect to {}: {}", endpoint_id, error));
}
return Err(anyhow::anyhow!("Connection to {} closed before it was established", endpoint_id));
}
None => {
return Err(anyhow::anyhow!("Connection stream ended before connecting to {}", endpoint_id));
}
}
}
}
}
}
pub async fn disconnect(&self, endpoint_id: iroh::EndpointId) -> anyhow::Result<()> {
self.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_DISCONNECTED_BY_USER,
)
.await
}
pub async fn is_current_transport_stable_id(
&self,
endpoint_id: iroh::EndpointId,
expected_transport_stable_id: u64,
) -> bool {
self.incoming_stream_generation_is_owned(endpoint_id, expected_transport_stable_id)
.await
}
pub async fn is_current_transport_stable_id_str(
&self,
endpoint_id: &str,
expected_transport_stable_id: u64,
) -> anyhow::Result<bool> {
let endpoint_id = endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|error| anyhow::anyhow!("invalid endpoint id: {error}"))?;
Ok(self
.is_current_transport_stable_id(endpoint_id, expected_transport_stable_id)
.await)
}
pub async fn disconnect_with_reason(
&self,
endpoint_id: iroh::EndpointId,
reason: &str,
) -> anyhow::Result<()> {
let endpoint_id_str = endpoint_id.to_string();
if std::env::var("PLUTO_RTC_TEARDOWN_TRACE").is_ok() {
eprintln!(
"[PlutoRTC][teardown-trace] Client::disconnect endpoint_id={} reason={}",
endpoint_id_str, reason
);
}
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
let records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
let transient_reconnect = crate::lifecycle_reason::is_transient(Some(reason));
node.disconnect_with_reason(endpoint_id, reason).await?;
#[cfg(any(feature = "transport-webrtc", feature = "transport-moq"))]
if !transient_reconnect {
node.clear_carrier_path_preference(endpoint_id);
}
for record in records {
if transient_reconnect {
self.connection_manager
.mark_transport_replaced(
&record.connection_id,
None,
Some("transient-disconnect".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string(),
),
)
.await;
#[cfg(target_arch = "wasm32")]
self.emit_current_wasm_connection_state(&record.connection_id)
.await;
#[cfg(not(target_arch = "wasm32"))]
self.emit_current_native_connection_state(&record.connection_id)
.await;
} else {
self.retire_managed_connection(&record.connection_id, Some(reason.to_string()))
.await;
}
}
#[cfg(not(target_arch = "wasm32"))]
if transient_reconnect {
let _ = self
.schedule_native_external_auto_connect_recovery(&endpoint_id_str)
.await;
}
Ok(())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
pub(crate) async fn disconnect_transport_generation_with_reason(
&self,
endpoint_id: iroh::EndpointId,
expected_transport_stable_id: u64,
reason: &str,
) -> anyhow::Result<bool> {
let endpoint_id_str = endpoint_id.to_string();
let records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
if !records
.iter()
.any(|record| record.transport_stable_id == Some(expected_transport_stable_id))
{
return Ok(false);
}
let node = self
.iroh_node
.read()
.await
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
if !node
.disconnect_with_reason_if_current(endpoint_id, expected_transport_stable_id, reason)
.await?
{
return Ok(false);
}
#[cfg(not(target_arch = "wasm32"))]
{
let mut replacement_observed = false;
for record in records {
replacement_observed |= self
.observe_native_transport_generation_lost(
&record.connection_id,
expected_transport_stable_id,
"liveness-probe-stale",
)
.await;
}
if replacement_observed {
let _ = self
.schedule_native_external_auto_connect_recovery(&endpoint_id.to_string())
.await;
}
}
#[cfg(target_arch = "wasm32")]
for record in records {
if self
.connection_manager
.mark_transport_replaced_if_current(
&record.connection_id,
expected_transport_stable_id,
Some("scoped-admission-fresh-redial".to_string()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string()),
)
.await
.is_some()
{
self.emit_current_wasm_connection_state(&record.connection_id)
.await;
}
}
Ok(true)
}
#[cfg(all(
not(target_arch = "wasm32"),
any(feature = "test-harness", feature = "test-relay-client")
))]
pub async fn test_disconnect(
&self,
endpoint_id: iroh::EndpointId,
transport_id: u64,
reason: &str,
) -> anyhow::Result<bool> {
self.disconnect_transport_generation_with_reason(endpoint_id, transport_id, reason)
.await
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn observe_native_transport_generation_lost(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
source: &str,
) -> bool {
let updated = self
.connection_manager
.mark_transport_replaced_if_current(
connection_id,
expected_transport_stable_id,
Some(source.to_string()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string()),
)
.await;
if updated.is_none() {
return false;
}
self.invalidate_native_main_route_proofs_for_transport(
connection_id,
expected_transport_stable_id,
);
true
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn observe_native_transport_generation_closed(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
close_reason: Option<&str>,
source: &str,
) -> bool {
if is_terminal_close_reason(close_reason) {
let manual_disconnect = is_manual_disconnect_close_reason(close_reason);
let reason = close_reason
.map(str::to_string)
.unwrap_or_else(|| crate::lifecycle_reason::REASON_CLOSED.to_string());
let closed = self
.connection_manager
.set_closed_if_current_or_replacement_handoff(
connection_id,
expected_transport_stable_id,
Some(reason),
)
.await;
if closed.is_none() {
return false;
}
self.invalidate_native_main_route_proofs_for_transport(
connection_id,
expected_transport_stable_id,
);
if manual_disconnect {
self.suppress_auto_connect_for_connection_peer(
connection_id,
"remote manual disconnect",
)
.await;
}
return true;
}
self.observe_native_transport_generation_lost(
connection_id,
expected_transport_stable_id,
source,
)
.await
}
pub fn policy_snapshot(&self) -> crate::runtime_policy::PolicySnapshot {
crate::runtime_policy::policy_snapshot()
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn open_bi(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
self.assert_raw_peer_stream_allowed(&endpoint_id).await?;
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
self.open_bi_internal(endpoint_id).await
}
pub(crate) async fn open_bi_internal(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
#[cfg(not(target_arch = "wasm32"))]
{
let (_, send, recv) = self
.open_bi_internal_with_transport_stable_id(endpoint_id)
.await?;
return Ok((send, recv));
}
#[cfg(target_arch = "wasm32")]
{
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
match node.open_bi(endpoint_id).await {
Ok(streams) => Ok(streams),
Err(first_error) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][iroh] open_bi cache miss/stale connection endpoint_id={} error={} action=redial",
endpoint_id, first_error
)));
self.ensure_connected(endpoint_id).await?;
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
node.open_bi(endpoint_id).await.map_err(|second_error| {
anyhow::anyhow!(
"open_bi failed after iroh redial endpoint_id={} first_error={} second_error={}",
endpoint_id,
first_error,
second_error
)
})
}
}
}
}
pub(crate) async fn open_current_bi_internal_with_transport_stable_id(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(u64, iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
node.open_bi_with_transport_stable_id(endpoint_id).await
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn open_bi_internal_with_transport_stable_id(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(u64, iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
match node.open_bi_with_transport_stable_id(endpoint_id).await {
Ok(streams) => Ok(streams),
Err(first_error) => {
eprintln!(
"[pluto-rtc][iroh] open_bi cache miss/stale connection endpoint_id={} error={} action=redial",
endpoint_id, first_error
);
self.ensure_connected_with_timeout(endpoint_id, std::time::Duration::from_secs(8))
.await?;
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
node.open_bi_with_transport_stable_id(endpoint_id)
.await
.map_err(|second_error| {
anyhow::anyhow!(
"open_bi failed after iroh redial endpoint_id={} first_error={} second_error={}",
endpoint_id,
first_error,
second_error
)
})
}
}
}
#[cfg(target_arch = "wasm32")]
pub(crate) async fn open_bi_internal_with_transport_stable_id(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(u64, iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
match node.open_bi_with_transport_stable_id(endpoint_id).await {
Ok(streams) => Ok(streams),
Err(first_error) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][iroh] open_bi cache miss/stale connection endpoint_id={} error={} action=redial",
endpoint_id, first_error
)));
self.ensure_connected(endpoint_id).await?;
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
node.open_bi_with_transport_stable_id(endpoint_id)
.await
.map_err(|second_error| {
anyhow::anyhow!(
"open_bi failed after iroh redial endpoint_id={} first_error={} second_error={}",
endpoint_id,
first_error,
second_error
)
})
}
}
}
pub(crate) async fn open_bi_internal_with_timeout(
&self,
endpoint_id: iroh::EndpointId,
timeout: Option<std::time::Duration>,
) -> anyhow::Result<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
#[cfg(not(target_arch = "wasm32"))]
{
if let Some(timeout) = timeout {
return tokio::time::timeout(timeout, self.open_bi_internal(endpoint_id))
.await
.map_err(|_| {
anyhow::anyhow!(
"open_bi timed out after {}ms endpoint_id={}",
timeout.as_millis(),
endpoint_id
)
})?;
}
}
#[cfg(target_arch = "wasm32")]
{
let _ = timeout;
}
self.open_bi_internal(endpoint_id).await
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn open_uni(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::SendStream> {
self.assert_raw_peer_stream_allowed(&endpoint_id).await?;
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
self.open_uni_internal(endpoint_id).await
}
pub(crate) async fn open_uni_internal(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::SendStream> {
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
match node.open_uni(endpoint_id).await {
Ok(stream) => Ok(stream),
Err(first_error) => {
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[pluto-rtc][iroh] open_uni cache miss/stale connection endpoint_id={} error={} action=redial",
endpoint_id, first_error
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][iroh] open_uni cache miss/stale connection endpoint_id={} error={} action=redial",
endpoint_id, first_error
)));
#[cfg(not(target_arch = "wasm32"))]
self.ensure_connected_with_timeout(endpoint_id, std::time::Duration::from_secs(8))
.await?;
#[cfg(target_arch = "wasm32")]
self.ensure_connected(endpoint_id).await?;
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
node.open_uni(endpoint_id).await.map_err(|second_error| {
anyhow::anyhow!(
"open_uni failed after iroh redial endpoint_id={} first_error={} second_error={}",
endpoint_id,
first_error,
second_error
)
})
}
}
}
#[cfg(target_arch = "wasm32")]
pub async fn open_uni(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::SendStream> {
self.ensure_connection_manager_record_before_peer_stream(&endpoint_id)
.await?;
self.open_uni_internal(endpoint_id).await
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn subscribe_accept_events(
&self,
) -> anyhow::Result<futures::stream::BoxStream<'static, AcceptEvent>> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
Ok(node.accept_events())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[cfg(target_arch = "wasm32")]
pub async fn subscribe_accept_events(
&self,
) -> anyhow::Result<n0_future::boxed::BoxStream<AcceptEvent>> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
Ok(node.accept_events())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn incoming_streams(
&self,
) -> anyhow::Result<async_channel::Receiver<IncomingStream>> {
let node_guard = self.iroh_node.read().await;
if node_guard.is_some() {
eprintln!(
"[pluto-rtc][native] incoming stream receiver attached node_id={}",
self.current_node_id().await.as_deref().unwrap_or("unknown")
);
Ok(self.native_application_streams_receiver.clone())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
pub async fn set_node_id(&self, node_id: String) {
let mut guard = self.node_id.write().await;
*guard = Some(node_id);
}
pub(crate) fn deterministic_connection_id(node_id_a: &str, node_id_b: &str) -> String {
if node_id_a <= node_id_b {
format!("{}-{}", node_id_a, node_id_b)
} else {
format!("{}-{}", node_id_b, node_id_a)
}
}
pub(crate) async fn known_remote_device_id_for_incoming_transport(
&self,
connection_id: &str,
remote_node_id: &str,
) -> Option<String> {
let from_record = |record: &crate::connection_manager::ConnectionRecord| {
record
.device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
record
.device_id_hint
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
})
};
if let Some(device_id) = self
.known_device_ids_by_node
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(remote_node_id)
.cloned()
{
return Some(device_id);
}
if let Some(existing) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
{
if let Some(device_id) = from_record(&existing) {
return Some(device_id);
}
}
let node_records = self.connection_manager.get_by_node_id(remote_node_id).await;
if let Some(device_id) = node_records
.iter()
.max_by(|left, right| left.updated_at_ms.cmp(&right.updated_at_ms))
.and_then(from_record)
{
return Some(device_id);
}
let auto_connect_user_id = {
let guard = match self.auto_connect_loop_key.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
guard.as_ref().map(|(user_id, _)| user_id.clone())
};
let Some(user_id) = auto_connect_user_id else {
return None;
};
self.search_devices(&user_id)
.await
.ok()
.and_then(|devices| {
devices.into_iter().find_map(|device| {
if device.node_id.as_deref() == Some(remote_node_id) {
Some(device.device_id)
} else {
None
}
})
})
}
pub(crate) async fn reconcile_authoritative_device_node(
&self,
remote_device_id: &str,
expected_node_id: &str,
) {
let remote_device_id = remote_device_id.trim();
let expected_node_id = expected_node_id.trim();
if remote_device_id.is_empty() || expected_node_id.is_empty() {
return;
}
self.observe_authoritative_device_node(remote_device_id, expected_node_id);
let records = self
.connection_manager
.get_by_device_id(remote_device_id)
.await;
let mut retired = 0usize;
let mut disconnected = 0usize;
let mut seen = std::collections::HashSet::new();
for record in records {
let Some(record_node_id) = record.node_id.clone() else {
continue;
};
if record_node_id == expected_node_id {
continue;
}
if !seen.insert(record.connection_id.clone()) {
continue;
}
if let Ok(endpoint_id) = record_node_id.parse::<iroh::EndpointId>() {
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_RETIRED_CONFLICTING_RECORD,
)
.await;
disconnected = disconnected.saturating_add(1);
}
self.retire_managed_connection_now(
&record.connection_id,
Some(crate::lifecycle_reason::REASON_RETIRED_CONFLICTING_RECORD.to_string()),
)
.await;
retired = retired.saturating_add(1);
}
if retired > 0 {
let msg = format!(
"[pluto-rtc][device-node-reconcile] retired conflicting records remote_device_id={} expected_node_id={} retired={} disconnected={}",
remote_device_id, expected_node_id, retired, disconnected
);
#[cfg(not(target_arch = "wasm32"))]
eprintln!("{}", msg);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&msg));
}
}
pub(crate) fn observe_authoritative_device_node(
&self,
remote_device_id: &str,
expected_node_id: &str,
) {
let remote_device_id = remote_device_id.trim();
let expected_node_id = expected_node_id.trim();
if remote_device_id.is_empty() || expected_node_id.is_empty() {
return;
}
let mut known = self
.known_device_ids_by_node
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
known.retain(|node_id, device_id| {
device_id != remote_device_id || node_id == expected_node_id
});
known.insert(expected_node_id.to_string(), remote_device_id.to_string());
drop(known);
self.known_device_endpoint_revision
.send_modify(|revision| *revision = revision.wrapping_add(1));
}
pub(crate) fn authoritative_node_for_device(&self, remote_device_id: &str) -> Option<String> {
let remote_device_id = remote_device_id.trim();
if remote_device_id.is_empty() {
return None;
}
self.known_device_ids_by_node
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.iter()
.find_map(|(node_id, device_id)| {
(device_id == remote_device_id).then(|| node_id.clone())
})
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn maybe_refresh_gateway_presence_for_auto_connect(
&self,
user_id: &str,
local_device_id: &str,
reason: &str,
last_republish_at_ms: &mut i64,
min_interval_ms: i64,
) {
let now = now_millis_i64();
if now.saturating_sub(*last_republish_at_ms) < min_interval_ms {
return;
}
*last_republish_at_ms = now;
let queued = self.request_presence_update();
eprintln!(
"[pluto-rtc][auto-connect][gateway-presence-refresh] user_id={} local_device_id={} reason={} owner=native-presence-actor queued={} liveness_source=gateway",
user_id, local_device_id, reason, queued
);
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn maybe_force_network_change_for_auto_connect(
&self,
user_id: &str,
local_device_id: &str,
reason: &str,
last_network_change_at_ms: &mut i64,
network_change_recovery_interval_ms: i64,
last_presence_republish_at_ms: &mut i64,
presence_republish_interval_ms: i64,
) {
let now = now_millis_i64();
if now.saturating_sub(*last_network_change_at_ms) < network_change_recovery_interval_ms {
return;
}
*last_network_change_at_ms = now;
match self.notify_network_change().await {
Ok(retired_stale) => {
eprintln!(
"[pluto-rtc][auto-connect] triggered network-change recovery user_id={} local_device_id={} reason={} retired_stale_records={}",
user_id, local_device_id, reason, retired_stale
);
}
Err(error) => {
eprintln!(
"[pluto-rtc][auto-connect] network-change recovery failed user_id={} local_device_id={} reason={} error={}",
user_id, local_device_id, reason, error
);
}
}
self.maybe_refresh_gateway_presence_for_auto_connect(
user_id,
local_device_id,
"network-change-recovery",
last_presence_republish_at_ms,
presence_republish_interval_ms,
)
.await;
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn handle_incoming_connection(
&self,
connection: iroh::endpoint::Connection,
) -> anyhow::Result<()> {
let remote_endpoint_id = connection.remote_id();
let remote_node_id = remote_endpoint_id.to_string();
let local_node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
let stable_id = crate::transport_generation::for_connection(&connection);
let node = {
let node_guard = self.iroh_node.read().await;
node_guard
.as_ref()
.cloned()
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?
};
let connection_for_close_check = connection.clone();
let (install_outcome_sender, install_outcome_receiver) = tokio::sync::oneshot::channel();
let accept_node = node.clone();
let accept_connection_for_task = connection.clone();
let mut accept_connection_task = tokio::spawn(async move {
accept_node
.accept_external_connection_with_install_notifier(
accept_connection_for_task,
install_outcome_sender,
)
.await
});
let install_outcome = tokio::select! {
biased;
outcome = install_outcome_receiver => outcome.map_err(|_| {
anyhow::anyhow!(
"native connection loop ended without an install decision for {}",
remote_node_id
)
})?,
result = &mut accept_connection_task => {
result.map_err(|error| {
anyhow::anyhow!(
"native connection loop task failed before install decision for {}: {}",
remote_node_id,
error,
)
})??;
anyhow::bail!(
"native connection loop ended before reporting an install decision for {}",
remote_node_id
);
}
};
match install_outcome {
crate::native_node::ExternalConnectionInstallOutcome::Installed {
transport_stable_id,
..
} if transport_stable_id == stable_id => {}
crate::native_node::ExternalConnectionInstallOutcome::Installed {
transport_stable_id,
..
} => {
connection.close(0u8.into(), b"native-install-generation-mismatch");
let _ = accept_connection_task.await;
anyhow::bail!(
"native install decision stable ID mismatch for {}: expected {}, got {}",
remote_node_id,
stable_id,
transport_stable_id
);
}
crate::native_node::ExternalConnectionInstallOutcome::KeptExisting {
fresh_transport_stable_id,
kept_transport_stable_id,
} => {
let result = accept_connection_task.await.map_err(|error| {
anyhow::anyhow!(
"native connection loop task failed for {}: {}",
remote_node_id,
error,
)
})?;
eprintln!(
"[PlutoRTC] handle_incoming_connection arbitration kept existing connection_id={} remote_node_id={} fresh_stable_id={} kept_stable_id={}",
connection_id,
remote_node_id,
fresh_transport_stable_id,
kept_transport_stable_id,
);
self.finalize_transport_dial_record(
remote_endpoint_id,
Some("incoming-kept-existing".to_string()),
)
.await?;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
return result;
}
crate::native_node::ExternalConnectionInstallOutcome::PendingReplacement {
transport_stable_id,
} => {
eprintln!(
"[PlutoRTC] handle_incoming_connection parked protected replacement connection_id={} remote_node_id={} transport_stable_id={}",
connection_id, remote_node_id, transport_stable_id,
);
drop(accept_connection_task);
return Ok(());
}
}
let prior_managed_snapshot = self.connection_manager.peer_snapshot(&connection_id).await;
let already_managed_connected = prior_managed_snapshot.as_ref().is_some_and(|snapshot| {
matches!(
snapshot.status,
crate::connection_manager::ConnectionState::Connected
)
});
self.finalize_transport_dial_record(remote_endpoint_id, Some("incoming".to_string()))
.await?;
let active_stable_id = self
.get_connection(remote_endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if active_stable_id != Some(stable_id) {
eprintln!(
"[PlutoRTC] handle_incoming_connection superseded before logical commit connection_id={} remote_node_id={} handler_stable_id={} active_stable_id={:?}",
connection_id, remote_node_id, stable_id, active_stable_id,
);
let result = accept_connection_task.await.map_err(|error| {
anyhow::anyhow!(
"native connection loop task failed for {}: {}",
remote_node_id,
error
)
})?;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
return result;
}
self.restore_scoped_route_admission(remote_endpoint_id, &connection_id)
.await?;
let _ = self
.maybe_start_offline_proof_for_transport(&connection_id, remote_endpoint_id, stable_id)
.await;
let replaced_main_route =
!self.native_admission_route_is_ready_for_transport(&connection_id, Some(stable_id));
let promoted_user_device = self
.promote_known_native_user_device_connection(&connection_id, &remote_node_id)
.await
.is_some();
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
let force_carrier_restart = incoming_transport_requires_carrier_restart(
already_managed_connected,
replaced_main_route,
);
eprintln!(
"[PlutoRTC] handle_incoming_connection registered connection_id={} local_node_id={} remote_node_id={}",
connection_id,
local_node_id,
remote_node_id
);
if self.session_registry_active() {
if matches!(
self.session_admission(&connection_id),
crate::session_token::SessionAdmission::Rejected { .. }
) {
eprintln!(
"[PlutoRTC] handle_incoming_connection resetting prior rejection for fresh transport connection_id={} remote_node_id={}",
connection_id,
remote_node_id
);
self.forget_session_connection(&connection_id);
}
if route_repair_needed(promoted_user_device, replaced_main_route) {
let _ = self
.wake_native_external_auto_connect_for_replaced_route(&remote_node_id)
.await;
}
let client = self.clone();
let connection_id_for_timeout = connection_id.clone();
let remote_node_id_for_timeout = remote_node_id.clone();
let remote_endpoint_id_for_timeout = remote_endpoint_id;
let timeout_connection = connection.clone();
let timeout_transport_stable_id =
crate::transport_generation::for_connection(&timeout_connection);
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(
crate::client::SESSION_ADMISSION_TIMEOUT_MS,
))
.await;
if !client.session_registry_active() {
return;
}
if client
.admission_timeout_still_owns_transport(
&connection_id_for_timeout,
remote_endpoint_id_for_timeout,
timeout_transport_stable_id,
)
.await
{
let Some(_retirement_guard) = client
.session_token_registry
.try_begin_admission_retirement(&connection_id_for_timeout)
else {
eprintln!(
"[PlutoRTC][session-admission][timeout-fenced] connection_id={} reason=response-writer-in-flight",
connection_id_for_timeout
);
return;
};
if !client
.admission_timeout_still_owns_transport(
&connection_id_for_timeout,
remote_endpoint_id_for_timeout,
timeout_transport_stable_id,
)
.await
{
return;
}
if !client
.admission_timeout_still_owns_transport(
&connection_id_for_timeout,
remote_endpoint_id_for_timeout,
timeout_transport_stable_id,
)
.await
{
return;
}
let reason = "session-admission-timeout";
timeout_connection.close(0u8.into(), reason.as_bytes());
client.forget_session_connection(&connection_id_for_timeout);
client
.retire_managed_connection_now(
&connection_id_for_timeout,
Some(reason.to_string()),
)
.await;
eprintln!(
"[PlutoRTC] handle_incoming_connection timed out pending admission connection_id={} remote_node_id={}",
connection_id_for_timeout,
remote_node_id_for_timeout
);
}
});
}
#[cfg(not(target_arch = "wasm32"))]
self.spawn_iroh_path_watcher(
&connection_id,
&remote_node_id,
&connection,
force_carrier_restart,
);
let result = accept_connection_task.await.map_err(|error| {
anyhow::anyhow!(
"native connection loop task failed for {}: {}",
remote_node_id,
error
)
})?;
match &result {
Ok(_) => {
let close_reason = connection_for_close_check.close_reason();
let manual_disconnect_notice = node
.take_manual_disconnect_notice(remote_endpoint_id, stable_id)
.await;
let close_reason_debug = format!("{:?}", close_reason);
eprintln!(
"[PlutoRTC] handle_incoming_connection stream loop exited connection_id={} close_reason={:?}",
connection_id,
close_reason
);
if manual_disconnect_notice {
let closed = self
.connection_manager
.set_closed_if_current_or_replacement_handoff(
&connection_id,
stable_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
if closed.is_none() {
eprintln!(
"[PlutoRTC][session-admission][manual-disconnect-fenced] ignoring stale notice connection_id={} remote_node_id={} closed_stable_id={}",
connection_id,
remote_node_id,
stable_id,
);
} else {
self.invalidate_native_main_route_proofs_for_transport(
&connection_id,
stable_id,
);
self.suppress_auto_connect_for_connection_peer(
&connection_id,
"remote manual disconnect notice",
)
.await;
}
} else if close_reason.is_none() {
if self
.connection_manager
.set_closed_if_current(&connection_id, stable_id, None)
.await
.is_some()
{
self.invalidate_native_main_route_proofs_for_transport(
&connection_id,
stable_id,
);
}
} else {
let resolution = self
.reconcile_incoming_transport_after_local_close(
&connection_id,
remote_endpoint_id,
&remote_node_id,
stable_id,
close_reason_debug.as_str(),
)
.await;
match resolution {
IncomingTransportCloseResolution::PreservedKeptTransport => {
eprintln!(
"[PlutoRTC] handle_incoming_connection: {} deduplicated (locally closed); \
preserving active manager record for kept connection",
connection_id,
);
}
IncomingTransportCloseResolution::RetiredClosedTransport => {
eprintln!(
"[PlutoRTC] handle_incoming_connection: {} deduplicated (locally closed); \
retired closed transport without a usable replacement",
connection_id,
);
}
}
}
}
Err(error) => {
eprintln!(
"[PlutoRTC] handle_incoming_connection failed connection_id={} error={}",
connection_id, error
);
let failed = self
.connection_manager
.set_failed_if_current(&connection_id, stable_id, Some(error.to_string()))
.await;
if failed.is_some() {
self.invalidate_native_main_route_proofs_for_transport(
&connection_id,
stable_id,
);
}
}
}
result
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn reconcile_incoming_transport_after_local_close(
&self,
connection_id: &str,
remote_endpoint_id: iroh::EndpointId,
remote_node_id: &str,
closed_stable_id: u64,
close_reason_debug: &str,
) -> IncomingTransportCloseResolution {
if is_manual_disconnect_close_reason(Some(close_reason_debug)) {
let closed = self
.connection_manager
.set_closed_if_current_or_replacement_handoff(
connection_id,
closed_stable_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
if closed.is_some() {
self.invalidate_native_main_route_proofs_for_transport(
connection_id,
closed_stable_id,
);
self.suppress_auto_connect_for_connection_peer(
connection_id,
"remote manual disconnect",
)
.await;
return IncomingTransportCloseResolution::RetiredClosedTransport;
}
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
if let Some(current) = self.get_connection(remote_endpoint_id).await {
let kept_stable_id = crate::transport_generation::for_connection(¤t);
let kept_transport_alive = current.close_reason().is_none();
let kept_transport_healthy = kept_transport_alive;
let handling = duplicate_closed_handling(
close_reason_debug,
kept_transport_alive,
kept_transport_healthy,
kept_stable_id,
closed_stable_id,
);
if handling == DuplicateClosedHandling::RebindKeptTransport {
eprintln!(
"[PlutoRTC] handle_incoming_connection rebinding manager record connection_id={} remote_node_id={} closed_stable_id={} kept_stable_id={} close_reason={} health_gate=preserve",
connection_id,
remote_node_id,
closed_stable_id,
kept_stable_id,
close_reason_debug,
);
let rebound = self
.connection_manager
.mark_transport_replaced_by_if_current(
connection_id,
closed_stable_id,
kept_stable_id,
Some("incoming-kept-existing".to_string()),
)
.await;
if rebound.is_some() {
if self.native_admission_route_is_ready_for_transport(
connection_id,
Some(kept_stable_id),
) {
let _ = self
.confirm_managed_connection_readiness_from_transport_proof(
connection_id,
kept_stable_id,
)
.await;
} else {
let _ = self
.confirm_managed_connection_readiness(connection_id)
.await;
}
}
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
eprintln!(
"[PlutoRTC] handle_incoming_connection duplicate-kept-existing health gate failed connection_id={} remote_node_id={} closed_stable_id={} kept_stable_id={} transport_alive={} transport_healthy={} close_reason={} handling={:?}",
connection_id,
remote_node_id,
closed_stable_id,
kept_stable_id,
kept_transport_alive,
kept_transport_healthy,
close_reason_debug,
handling,
);
}
if is_replacement_churn_close_reason(Some(close_reason_debug)) {
eprintln!(
"[PlutoRTC] handle_incoming_connection preserving replacement-churn close connection_id={} remote_node_id={} closed_stable_id={} close_reason={}",
connection_id,
remote_node_id,
closed_stable_id,
close_reason_debug,
);
self.connection_manager
.mark_transport_replaced_if_current(
connection_id,
closed_stable_id,
Some("incoming-replacement-churn".to_string()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string()),
)
.await;
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
if !is_terminal_close_reason(Some(close_reason_debug)) {
let recovery_source = format!("incoming-transport-lost:{close_reason_debug}");
let replacement_observed = self
.observe_native_transport_generation_lost(
connection_id,
closed_stable_id,
&recovery_source,
)
.await
|| self
.connection_manager
.replacement_handoff_matches(connection_id, closed_stable_id)
.await;
if replacement_observed {
eprintln!(
"[PlutoRTC] handle_incoming_connection preserving recoverable transport loss connection_id={} remote_node_id={} closed_stable_id={} close_reason={}",
connection_id,
remote_node_id,
closed_stable_id,
close_reason_debug,
);
let _ = self
.schedule_native_external_auto_connect_recovery(&remote_endpoint_id.to_string())
.await;
}
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
eprintln!(
"[PlutoRTC] handle_incoming_connection retiring stale manager record connection_id={} remote_node_id={} closed_stable_id={} close_reason={}",
connection_id,
remote_node_id,
closed_stable_id,
close_reason_debug,
);
let terminal_reason = crate::lifecycle_reason::LifecycleReasonCode::from_text(Some(
close_reason_debug,
))
.filter(|reason| reason.is_terminal())
.map(crate::lifecycle_reason::LifecycleReasonCode::as_str)
.unwrap_or(
crate::lifecycle_reason::REASON_INCOMING_TRANSPORT_CLOSED_WITHOUT_LIVE_REPLACEMENT,
);
let closed = self
.connection_manager
.set_closed_if_current_or_replacement_handoff(
connection_id,
closed_stable_id,
Some(terminal_reason.to_string()),
)
.await;
if closed.is_some() {
self.invalidate_native_main_route_proofs_for_transport(connection_id, closed_stable_id);
IncomingTransportCloseResolution::RetiredClosedTransport
} else {
IncomingTransportCloseResolution::PreservedKeptTransport
}
}
async fn suppress_auto_connect_for_connection_peer(&self, connection_id: &str, reason: &str) {
let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return;
};
let device_id = record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string);
let device_id = match device_id {
Some(device_id) => device_id,
None => {
let Some(remote_node_id) = record
.node_id
.as_deref()
.or(record.endpoint_id.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty())
else {
return;
};
let Some((user_id, _)) = self.active_session_identity() else {
return;
};
let Ok(devices) = self.search_devices(&user_id).await else {
return;
};
let Some(device) = devices.into_iter().find(|device| {
device
.node_id
.as_deref()
.map(str::trim)
.is_some_and(|node_id| node_id.eq_ignore_ascii_case(remote_node_id))
}) else {
return;
};
let device_id = device.device_id.trim().to_string();
if device_id.is_empty() {
return;
}
let _ = self
.connection_manager
.set_device_id(connection_id, device_id.clone())
.await;
device_id
}
};
let node_alias = record
.node_id
.as_deref()
.or(record.endpoint_id.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty());
self.set_peer_requested_auto_connect_excluded(&device_id, node_alias);
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[PlutoRTC] auto-connect suppressed for peer device_id={} connection_id={} reason={}",
device_id, connection_id, reason
);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc] auto-connect suppressed for peer device_id={} connection_id={} reason={}",
device_id, connection_id, reason
)));
}
async fn send_manual_disconnect_notice(&self, endpoint_id: iroh::EndpointId) {
let frame = crate::heartbeat::codec::disconnect_frame();
let Ok(send) = self.open_uni(endpoint_id).await else {
return;
};
let mut send = crate::application_crypto_streams::PeerSendStream::plain(send);
if send.write_all(&frame).await.is_ok() {
let _ = send
.finish_and_wait_for_peer(std::time::Duration::from_secs(2))
.await;
}
}
pub async fn disconnect_device(
self: &std::sync::Arc<Self>,
device_id: &str,
node_id_hint: Option<&str>,
) -> Vec<String> {
let device_id_trim = device_id.trim();
if !device_id_trim.is_empty() {
self.exclude_peer_and_publish(device_id_trim).await;
if let Some(node_id) = node_id_hint
.map(str::trim)
.filter(|value| !value.is_empty())
{
self.set_auto_connect_excluded_peer(device_id_trim, Some(node_id), true);
}
}
let mut records = if !device_id_trim.is_empty() {
self.resolve_peer_connection_records(device_id_trim).await
} else {
Vec::new()
};
if records.is_empty() {
if let Some(node_id) = node_id_hint
.map(str::trim)
.filter(|value| !value.is_empty())
{
records = self.resolve_peer_connection_records(node_id).await;
}
}
if records.is_empty() && !device_id_trim.is_empty() {
if let Some((user_id, _)) = self.active_session_identity() {
if let Ok(devices) = self.search_devices(&user_id).await {
if let Some(device) = devices
.iter()
.find(|device| device.device_id.trim() == device_id_trim)
{
if let Some(node_id) = device
.node_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
self.set_auto_connect_excluded_peer(
device_id_trim,
Some(node_id),
true,
);
records = self.resolve_peer_connection_records(node_id).await;
}
}
}
}
}
let mut retired = Vec::with_capacity(records.len());
for record in records {
if !device_id_trim.is_empty() {
let node_alias = record
.node_id
.as_deref()
.or(record.endpoint_id.as_deref())
.map(str::trim)
.filter(|value| !value.is_empty());
self.set_auto_connect_excluded_peer(device_id_trim, node_alias, true);
}
let endpoint_id = record
.endpoint_id
.as_deref()
.or(record.node_id.as_deref())
.and_then(|value| value.parse::<iroh::EndpointId>().ok());
if let Some(endpoint_id) = endpoint_id {
self.send_manual_disconnect_notice(endpoint_id).await;
if let Err(error) = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_MANUAL_DISCONNECT,
)
.await
{
#[cfg(not(target_arch = "wasm32"))]
eprintln!(
"[PlutoRTC] disconnect_device close error connection_id={} endpoint_id={} error={}",
record.connection_id, endpoint_id, error
);
#[cfg(target_arch = "wasm32")]
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc] disconnect_device close error connection_id={} endpoint_id={} error={}",
record.connection_id, endpoint_id, error
)));
}
} else {
self.retire_managed_connection(
&record.connection_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
}
retired.push(record.connection_id);
}
retired
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn get_connection(
&self,
endpoint_id: iroh::EndpointId,
) -> Option<iroh::endpoint::Connection> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.get_connection(endpoint_id).await
} else {
None
}
}
#[cfg(target_arch = "wasm32")]
pub async fn get_connection(
&self,
endpoint_id: iroh::EndpointId,
) -> Option<iroh::endpoint::Connection> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.get_connection(endpoint_id).await
} else {
None
}
}
pub(super) async fn resolve_iroh_endpoint_id_for_peer(
&self,
peer_id: &str,
) -> Option<iroh::EndpointId> {
let peer_id = peer_id.trim();
if peer_id.is_empty() {
return None;
}
if let Ok(endpoint_id) = peer_id.parse::<iroh::EndpointId>() {
return Some(endpoint_id);
}
if let Some(record) = self.connection_manager.get_by_connection_id(peer_id).await {
for raw in [record.endpoint_id.as_deref(), record.node_id.as_deref()]
.into_iter()
.flatten()
{
if let Ok(endpoint_id) = raw.trim().parse::<iroh::EndpointId>() {
return Some(endpoint_id);
}
}
}
let snapshot = self.peer_session(peer_id).await?;
let node_id = snapshot.node_id.as_deref()?.trim();
node_id.parse::<iroh::EndpointId>().ok()
}
pub async fn iroh_transport_rtt_ms(&self, peer_id: &str) -> Option<u64> {
let endpoint_id = self.resolve_iroh_endpoint_id_for_peer(peer_id).await?;
let connection = self.get_connection(endpoint_id).await?;
let paths = connection.paths();
let rtt = paths
.iter()
.filter(|path| path.is_selected())
.map(|path| path.rtt())
.min()
.or_else(|| paths.iter().map(|path| path.rtt()).min())
.or_else(|| connection.rtt(iroh::endpoint::PathId::ZERO));
rtt.map(|value| value.as_millis().max(1) as u64)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn iroh_path_kind(&self, peer_id: &str) -> IrohPathKind {
let endpoint_id = match self.resolve_iroh_endpoint_id_for_peer(peer_id).await {
Some(id) => id,
None => return IrohPathKind::Unknown,
};
let connection = match self.get_connection(endpoint_id).await {
Some(c) => c,
None => return IrohPathKind::Unknown,
};
self.classify_native_iroh_connection_path(&connection).await
}
#[cfg(not(target_arch = "wasm32"))]
async fn classify_native_iroh_connection_path(
&self,
connection: &iroh::endpoint::Connection,
) -> IrohPathKind {
let custom_kinds = self.native_custom_transport_kinds.read().await;
for path in connection.paths().iter() {
if !path.is_selected() {
continue;
}
if let iroh::TransportAddr::Custom(addr) = path.remote_addr() {
if let Some(kind) = custom_kinds.get(&addr.id()) {
return *kind;
}
}
}
drop(custom_kinds);
#[cfg(feature = "transport-lan")]
{
return crate::local_discovery::classify_iroh_path_kind(&connection);
}
#[cfg(not(feature = "transport-lan"))]
{
let paths = connection.paths();
let has_direct = paths.iter().any(|p| p.is_selected() && p.is_ip());
let has_relay = paths.iter().any(|p| p.is_selected() && p.is_relay());
if has_direct {
IrohPathKind::DirectQuic
} else if has_relay {
IrohPathKind::Relay
} else {
IrohPathKind::Unknown
}
}
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
pub async fn list_local_peers(&self) -> Vec<crate::local_discovery::LocalPeerSnapshot> {
self.local_discovery_registry.list_peers().await
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
pub async fn is_peer_locally_reachable(&self, node_id: &str) -> bool {
self.local_discovery_registry
.is_locally_reachable(node_id)
.await
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
async fn lan_discovery_enabled(&self) -> bool {
if self.relay_only_mode_enabled().await {
return false;
}
let config = self.transport_config.read().await;
config
.iroh_lan
.as_ref()
.map(|lan| lan.enabled)
.unwrap_or(false)
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
async fn lan_discovery_advertise(&self) -> bool {
let config = self.transport_config.read().await;
config
.iroh_lan
.as_ref()
.map(|lan| lan.advertise)
.unwrap_or(true)
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
async fn apply_optional_lan_discovery(
&self,
builder: iroh::endpoint::Builder,
) -> anyhow::Result<iroh::endpoint::Builder> {
Ok(builder)
}
#[cfg(all(not(target_arch = "wasm32"), not(feature = "transport-lan")))]
async fn apply_optional_lan_discovery(
&self,
builder: iroh::endpoint::Builder,
) -> anyhow::Result<iroh::endpoint::Builder> {
Ok(builder)
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
async fn start_mdns_discovery_tasks(&self, endpoint: &iroh::Endpoint, node_id: &str) {
if !self.lan_discovery_enabled().await {
return;
}
let advertise = self.lan_discovery_advertise().await;
let mdns = match iroh_mdns_address_lookup::MdnsAddressLookup::builder()
.advertise(advertise)
.build(endpoint.id())
{
Ok(mdns) => mdns,
Err(error) => {
eprintln!(
"[pluto-rtc][lan] mdns build failed node_id={} error={}",
node_id, error
);
return;
}
};
if let Ok(lookup_services) = endpoint.address_lookup() {
lookup_services.add(mdns.clone());
}
*self.mdns_address_lookup.write().await = Some(mdns.clone());
crate::local_discovery::spawn_mdns_discovery_task(
mdns,
self.local_discovery_registry.clone(),
node_id.to_string(),
self.clone(),
);
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-lan"))]
async fn start_local_discovery_tasks(&self, endpoint: &iroh::Endpoint, node_id: &str) {
self.start_mdns_discovery_tasks(endpoint, node_id).await;
}
#[cfg(all(not(target_arch = "wasm32"), not(feature = "transport-lan")))]
async fn start_local_discovery_tasks(&self, _endpoint: &iroh::Endpoint, _node_id: &str) {}
#[cfg(target_arch = "wasm32")]
pub async fn iroh_path_kind(&self, peer_id: &str) -> IrohPathKind {
let Some(endpoint_id) = self.resolve_iroh_endpoint_id_for_peer(peer_id).await else {
return IrohPathKind::Unknown;
};
let Some(connection) = self.get_connection(endpoint_id).await else {
return IrohPathKind::Unknown;
};
for path in connection.paths().iter().filter(|path| path.is_selected()) {
let iroh::TransportAddr::Custom(addr) = path.remote_addr() else {
continue;
};
#[cfg(not(any(feature = "transport-webrtc", feature = "transport-moq")))]
let _ = addr;
#[cfg(feature = "transport-webrtc")]
if addr.id() == crate::iroh_carrier_kind::EXPERIMENTAL_WEBRTC_TRANSPORT_ID {
return IrohPathKind::WebRtc;
}
#[cfg(feature = "transport-moq")]
if addr.id() == crate::iroh_carrier_kind::EXPERIMENTAL_MOQ_TRANSPORT_ID {
return IrohPathKind::Moq;
}
}
IrohPathKind::Relay
}
pub(crate) async fn incoming_stream_generation_is_owned(
&self,
endpoint_id: iroh::EndpointId,
incoming_transport_stable_id: u64,
) -> bool {
self.incoming_stream_generation_fence(endpoint_id, incoming_transport_stable_id)
.await
.owned
}
#[cfg(target_arch = "wasm32")]
pub(crate) async fn incoming_stream_generation_is_authenticated_reconnect_candidate(
&self,
endpoint_id: iroh::EndpointId,
incoming_transport_stable_id: u64,
) -> bool {
let fence = self
.incoming_stream_generation_fence(endpoint_id, incoming_transport_stable_id)
.await;
if fence.owned {
return false;
}
let Some(local_node_id) = self.current_node_id().await else {
return false;
};
let remote_node_id = endpoint_id.to_string();
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
let record = self
.connection_manager
.get_by_connection_id(&connection_id)
.await;
let device_id = record.as_ref().and_then(|record| {
record
.device_id
.as_deref()
.or(record.device_id_hint.as_deref())
});
let peer_requested_exclusion =
self.auto_connect_exclusion_for_peer(device_id, Some(&remote_node_id)) == Some(true);
self.connection_manager
.incoming_authenticated_reconnect_candidate(
&connection_id,
fence.physical_transport_stable_id,
incoming_transport_stable_id,
peer_requested_exclusion,
)
.await
}
pub(crate) async fn incoming_stream_generation_fence(
&self,
endpoint_id: iroh::EndpointId,
incoming_transport_stable_id: u64,
) -> crate::connection_manager::IncomingTransportGenerationFence {
let physical_transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let Some(local_node_id) = self.current_node_id().await else {
return crate::connection_manager::IncomingTransportGenerationFence {
incoming_transport_stable_id,
physical_transport_stable_id,
managed_transport_stable_id: None,
manager_record_exists: false,
owned: false,
};
};
let connection_id =
Self::deterministic_connection_id(&local_node_id, &endpoint_id.to_string());
self.connection_manager
.incoming_transport_generation_fence(
&connection_id,
physical_transport_stable_id,
incoming_transport_stable_id,
)
.await
}
#[cfg(target_arch = "wasm32")]
async fn promote_wasm_incoming_transport(
&self,
endpoint_id: iroh::EndpointId,
source: &str,
) -> anyhow::Result<()> {
let remote_node_id = endpoint_id.to_string();
let local_node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let connection_id = Self::deterministic_connection_id(&local_node_id, &remote_node_id);
let stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection))
.ok_or_else(|| {
anyhow::anyhow!(
"cannot promote WASM base transport without a physical stable id for {}",
remote_node_id
)
})?;
let desired_device_id = self
.current_desired_user_device_for_node(&remote_node_id)
.await;
let existing = self
.connection_manager
.get_by_connection_id(&connection_id)
.await;
let needs_desired_peer_reopen = existing.as_ref().is_none_or(|record| {
matches!(
record.state,
crate::connection_manager::ConnectionState::Closing
| crate::connection_manager::ConnectionState::Closed
| crate::connection_manager::ConnectionState::Failed
)
});
if let Some(device_id) = desired_user_device_reopen_candidate(
needs_desired_peer_reopen,
desired_device_id.as_deref(),
) {
if self.is_auto_connect_excluded(device_id) {
return Err(anyhow::anyhow!(
"WASM incoming transport is excluded by local peer intent"
));
}
self.connection_manager
.begin_automatic_connect(
connection_id.clone(),
Some(remote_node_id.clone()),
Some(device_id.to_string()),
Some(remote_node_id.clone()),
|_| true,
)
.await
.ok_or_else(|| {
anyhow::anyhow!("WASM incoming transport was fenced by terminal peer intent")
})?;
if self
.current_desired_user_device_for_node(&remote_node_id)
.await
.as_deref()
!= Some(device_id)
{
return Err(anyhow::anyhow!(
"WASM incoming transport was superseded by newer desired-peer evidence"
));
}
}
let already_promoted = self
.connection_manager
.get_by_connection_id(&connection_id)
.await
.map(|record| {
record.state == crate::connection_manager::ConnectionState::Connected
&& record.transport_stable_id == Some(stable_id)
})
.unwrap_or(false);
let committed = match desired_device_id.as_deref() {
Some(device_id) => {
self.commit_current_desired_user_device_transport_record(
&connection_id,
&remote_node_id,
device_id,
stable_id,
Some(source.to_string()),
)
.await
}
None => {
self.commit_current_native_transport_record(
&connection_id,
&remote_node_id,
None,
Some(stable_id),
Some(source.to_string()),
)
.await
}
};
if !committed {
let current_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
let current_desired_device_id = self
.current_desired_user_device_for_node(&remote_node_id)
.await;
let denial = if current_stable_id != Some(stable_id) {
"physical-generation-stale"
} else if desired_device_id.is_none() {
"desired-assignment-absent"
} else if current_desired_device_id != desired_device_id {
"desired-assignment-stale"
} else {
"logical-or-security-denied"
};
return Err(anyhow::anyhow!(
"WASM physical transport commit rejected reason={denial} stable_id={stable_id} current_stable_id={current_stable_id:?}"
));
}
self.restore_scoped_route_admission(endpoint_id, &connection_id)
.await?;
let active_transport = self.iroh_path_kind(&remote_node_id).await.transport_label();
let _ = self
.report_transport_status_for_current_generation(&connection_id, active_transport, None)
.await;
if let crate::session_token::SessionAdmission::Accepted {
scope: Some(scope), ..
} = self.session_token_registry.admission(&connection_id)
{
let scope_name = scope.as_str().trim();
if !scope_name.is_empty() {
self.connection_manager
.add_scope(&connection_id, scope_name)
.await;
}
}
if let Some(device_id) = self
.known_remote_device_id_for_incoming_transport(&connection_id, &remote_node_id)
.await
{
self.mark_trusted_user_device_connection_admitted(&connection_id, &device_id)
.await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] bound incoming transport connection_id={} remote_node_id={} device_id={} source={}",
connection_id, remote_node_id, device_id, source
)));
}
self.emit_current_wasm_connection_state(&connection_id)
.await;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
if !already_promoted {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] promoted incoming transport connection_id={} remote_node_id={} stable_id={} source={}",
connection_id, remote_node_id, stable_id, source
)));
}
Ok(())
}
#[cfg(target_arch = "wasm32")]
async fn handle_wasm_accept_event(self: &Arc<Self>, event: AcceptEvent) {
match event {
AcceptEvent::Accepted {
endpoint_id,
transport_stable_id,
} => {
let current_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if !accepted_transport_event_is_current(transport_stable_id, current_stable_id) {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] ignored stale accepted transport endpoint_id={} accepted_stable_id={} current_stable_id={:?}",
endpoint_id, transport_stable_id, current_stable_id
)));
return;
}
if let Some(withdrawn_device_id) = self
.withdrawn_desired_device_for_node(&endpoint_id.to_string())
.await
{
let _ = self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_DESIRED_PEER_WITHDRAWN,
)
.await;
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] rejected withdrawn desired peer endpoint_id={} device_id={} transport_stable_id={}",
endpoint_id, withdrawn_device_id, transport_stable_id,
)));
return;
}
if let Err(error) = self
.promote_wasm_incoming_transport(endpoint_id, "wasm-accept")
.await
{
let node = self.iroh_node.read().await.as_ref().cloned();
let rejected_closed = match node {
Some(node) => node
.disconnect_with_reason_if_current(
endpoint_id,
transport_stable_id,
crate::lifecycle_reason::REASON_AUTO_CONNECT_FAILURE,
)
.await
.unwrap_or(false),
None => false,
};
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] failed promoting accepted transport endpoint_id={} transport_stable_id={} rejected_closed={} error={:#}",
endpoint_id, transport_stable_id, rejected_closed, error
)));
self.wake_browser_auto_connect();
}
}
AcceptEvent::Closed {
endpoint_id,
transport_stable_id,
error,
was_locally_closed,
} => {
let remote_node_id = endpoint_id.to_string();
let Some(local_node_id) = self.current_node_id().await else {
return;
};
let connection_id =
Self::deterministic_connection_id(&local_node_id, &remote_node_id);
if !self
.connection_manager
.current_transport_matches(&connection_id, Some(transport_stable_id))
.await
{
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] ignored stale accepted-transport close endpoint_id={} connection_id={} closed_stable_id={} error={:?}",
remote_node_id, connection_id, transport_stable_id, error
)));
return;
}
if !self.is_connected(endpoint_id).await {
if is_terminal_close_reason(error.as_deref()) {
let manual_disconnect = is_manual_disconnect_close_reason(error.as_deref());
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] terminal remote close — closing without replacement wait endpoint_id={} connection_id={} error={:?}",
remote_node_id, connection_id, error
)));
let closed = self
.connection_manager
.set_closed_if_current(
&connection_id,
transport_stable_id,
Some(error.unwrap_or_else(|| {
crate::lifecycle_reason::REASON_CLOSED.to_string()
})),
)
.await;
if closed.is_none() {
return;
}
if manual_disconnect {
self.suppress_auto_connect_for_connection_peer(
&connection_id,
"remote manual disconnect",
)
.await;
}
self.emit_current_wasm_connection_state(&connection_id)
.await;
return;
}
if was_locally_closed {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] locally-closed transport — skipping replacement wait endpoint_id={} connection_id={} error={:?}",
remote_node_id, connection_id, error
)));
self.connection_manager
.mark_transport_replaced_if_current(
&connection_id,
transport_stable_id,
Some("wasm-accept-locally-closed".to_string()),
Some("locally-closed".to_string()),
)
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
return;
}
if let Some(rebound_stable_id) = self
.await_wasm_replacement_transport(
endpoint_id,
"wasm-accept-close-grace",
crate::runtime_policy::WASM_ACCEPT_CLOSE_GRACE_MS,
)
.await
{
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] accepted transport close settled onto replacement endpoint_id={} connection_id={} rebound_stable_id={:?}",
remote_node_id,
connection_id,
rebound_stable_id
)));
return;
}
if is_replacement_churn_close_reason(error.as_deref()) {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] replacement-churn close awaiting successor endpoint_id={} connection_id={} closed_stable_id={} error={:?}",
remote_node_id, connection_id, transport_stable_id, error
)));
self.connection_manager
.mark_transport_replaced_if_current(
&connection_id,
transport_stable_id,
Some("wasm-accept-replacement-churn".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS
.to_string(),
),
)
.await;
} else {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] accepted transport closed without live replacement endpoint_id={} connection_id={} closed_stable_id={} error={:?}",
remote_node_id, connection_id, transport_stable_id, error
)));
self.connection_manager
.set_closed_if_current_or_replacement_handoff(
&connection_id,
transport_stable_id,
Some(
crate::lifecycle_reason::REASON_INCOMING_TRANSPORT_CLOSED_WITHOUT_LIVE_REPLACEMENT
.to_string(),
),
)
.await;
}
self.emit_current_wasm_connection_state(&connection_id)
.await;
self.wake_browser_auto_connect();
} else {
let new_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if let Some(new_stable_id) = new_stable_id {
self.connection_manager
.mark_transport_replaced_by_if_current(
&connection_id,
transport_stable_id,
new_stable_id,
Some("wasm-accept-close-rebind".to_string()),
)
.await;
}
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] accepted transport closed but live replacement exists endpoint_id={} connection_id={} rebound_stable_id={:?}",
remote_node_id,
connection_id,
new_stable_id
)));
self.emit_current_wasm_connection_state(&connection_id)
.await;
}
}
}
}
#[cfg(target_arch = "wasm32")]
pub fn start_wasm_accept_bridge(self: Arc<Self>) {
if self
.wasm_accept_bridge_started
.compare_exchange(
false,
true,
std::sync::atomic::Ordering::SeqCst,
std::sync::atomic::Ordering::SeqCst,
)
.is_err()
{
return;
}
n0_future::task::spawn(async move {
let mut events = match self.subscribe_accept_events().await {
Ok(events) => events,
Err(error) => {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] failed subscribing to accept events: {}",
error
)));
self.wasm_accept_bridge_started
.store(false, std::sync::atomic::Ordering::SeqCst);
return;
}
};
use futures::StreamExt;
while let Some(event) = events.next().await {
self.handle_wasm_accept_event(event).await;
}
self.wasm_accept_bridge_started
.store(false, std::sync::atomic::Ordering::SeqCst);
});
}
#[cfg(target_arch = "wasm32")]
async fn await_wasm_replacement_transport(
self: &Arc<Self>,
endpoint_id: iroh::EndpointId,
source: &str,
timeout_ms: u64,
) -> Option<u64> {
let started_ms = js_sys::Date::now();
let mut poll_interval_ms: u64 = 50;
const MAX_POLL_INTERVAL_MS: u64 =
crate::runtime_policy::WASM_REPLACEMENT_POLL_MAX_INTERVAL_MS;
loop {
let rebound_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| crate::transport_generation::for_connection(&connection));
if let Some(rebound_stable_id) = rebound_stable_id {
if self
.promote_wasm_incoming_transport(endpoint_id, source)
.await
.is_ok()
{
return Some(rebound_stable_id);
}
}
if (js_sys::Date::now() - started_ms) >= timeout_ms as f64 {
return None;
}
gloo_timers::future::sleep(std::time::Duration::from_millis(poll_interval_ms)).await;
poll_interval_ms = (poll_interval_ms * 2).min(MAX_POLL_INTERVAL_MS);
}
}
#[cfg(target_arch = "wasm32")]
fn start_wasm_connect_event_bridge<S>(
self: Arc<Self>,
endpoint_id: iroh::EndpointId,
closed_transport_stable_id: Option<u64>,
mut events: S,
) where
S: futures::Stream<Item = ConnectEvent> + Unpin + 'static,
{
n0_future::task::spawn(async move {
use futures::StreamExt;
while let Some(event) = events.next().await {
match event {
ConnectEvent::Connected => {}
ConnectEvent::Closed { error } => {
let endpoint_id_str = endpoint_id.to_string();
let duplicate_close =
is_duplicate_kept_existing_close_reason(error.as_deref());
let replacement_churn_close =
is_replacement_churn_close_reason(error.as_deref());
let terminal_disconnect_close = is_terminal_close_reason(error.as_deref());
let live_transport = self.is_connected(endpoint_id).await;
let current_live_stable_id =
self.get_connection(endpoint_id).await.map(|connection| {
crate::transport_generation::for_connection(&connection)
});
let connected_managed_stable_id = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await
.into_iter()
.find(|record| {
record.state
== crate::connection_manager::ConnectionState::Connected
&& record.transport_stable_id == current_live_stable_id
})
.and_then(|record| record.transport_stable_id);
match classify_wasm_closed_transport_rebind(
terminal_disconnect_close,
current_live_stable_id,
connected_managed_stable_id,
live_transport,
) {
WasmClosedTransportRebind::TerminalDisconnect => {
}
WasmClosedTransportRebind::AlreadyPromoted(stable_id) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] ignored stale close after current generation promotion endpoint_id={} closed_transport_stable_id={:?} current_live_stable_id={} error={:?}",
endpoint_id_str,
closed_transport_stable_id,
stable_id,
error
)));
continue;
}
WasmClosedTransportRebind::BaseReplacement(stable_id) => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] observed closed outgoing transport with concrete base replacement endpoint_id={} closed_transport_stable_id={:?} current_live_stable_id={} error={:?}",
endpoint_id_str,
closed_transport_stable_id,
stable_id,
error
)));
let _ = self
.promote_wasm_incoming_transport(
endpoint_id,
"wasm-connect-bridge-rebind",
)
.await;
continue;
}
WasmClosedTransportRebind::SurvivingCarrierAwaitingReplacement => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] committed transport closed while a carrier observation remains; awaiting concrete replacement endpoint_id={} closed_transport_stable_id={:?} error={:?}",
endpoint_id_str,
closed_transport_stable_id,
error
)));
continue;
}
WasmClosedTransportRebind::NoLiveTransport => {}
}
if duplicate_close || replacement_churn_close {
let dup_records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
for dup_record in &dup_records {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] replacement close: marking replacement-pending for incoming-accept hold endpoint_id={} connection_id={} error={:?}",
endpoint_id_str,
dup_record.connection_id,
error
)));
if let Some(closed_stable_id) = closed_transport_stable_id {
if self.connection_manager
.mark_transport_replaced_if_current(
&dup_record.connection_id,
closed_stable_id,
Some("replacement-held-for-accept".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS
.to_string(),
),
)
.await
.is_some()
{
self.emit_current_wasm_connection_state(&dup_record.connection_id)
.await;
}
}
}
if let Some(rebound_stable_id) = self
.await_wasm_replacement_transport(
endpoint_id,
"wasm-connect-duplicate-close-grace",
crate::runtime_policy::WASM_CONNECT_REPLACEMENT_GRACE_MS,
)
.await
{
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] replacement close settled onto replacement endpoint_id={} closed_transport_stable_id={:?} rebound_stable_id={:?}",
endpoint_id_str,
closed_transport_stable_id,
rebound_stable_id
)));
continue;
}
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] replacement close grace expired without replacement endpoint_id={} closed_transport_stable_id={:?} error={:?}",
endpoint_id_str,
closed_transport_stable_id,
error
)));
let reset = match closed_transport_stable_id {
Some(closed_stable_id) => {
let node = self.iroh_node.read().await.as_ref().cloned();
match node {
Some(node) => node
.disconnect_with_reason_if_current(
endpoint_id,
closed_stable_id,
crate::lifecycle_reason::REASON_ENDPOINT_HARD_RESET_REDIAL,
)
.await
.map(|disconnected| disconnected.then_some(())),
None => Ok(None),
}
}
None => Ok(None),
};
match reset {
Ok(Some(())) => {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] replacement close fallback reset endpoint_id={} closed_transport_stable_id={:?}",
endpoint_id_str,
closed_transport_stable_id
)));
}
Ok(None) => {}
Err(error) => {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] duplicate close fallback reset failed endpoint_id={} error={}",
endpoint_id_str,
error
)));
}
}
continue;
}
let records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] closed outgoing transport lookup endpoint_id={} records_found={} closed_transport_stable_id={:?} error={:?}",
endpoint_id_str,
records.len(),
closed_transport_stable_id,
error
)));
let mut stale_close_ignored = 0usize;
let mut records_to_retire = Vec::new();
for record in records {
let matches_closed_transport = match closed_transport_stable_id {
Some(closed_stable_id) => {
self.connection_manager
.current_transport_matches(
&record.connection_id,
Some(closed_stable_id),
)
.await
}
None => false,
};
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] record transport check connection_id={} record_stable_id={:?} closed_stable_id={:?} matches={}",
record.connection_id,
record.transport_stable_id,
closed_transport_stable_id,
matches_closed_transport
)));
if matches_closed_transport {
records_to_retire.push(record);
} else {
stale_close_ignored = stale_close_ignored.saturating_add(1);
}
}
if stale_close_ignored > 0 {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] ignored stale outgoing close endpoint_id={} ignored_records={} closed_transport_stable_id={:?}",
endpoint_id_str,
stale_close_ignored,
closed_transport_stable_id
)));
}
let retire_count = records_to_retire.len();
for record in records_to_retire {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] evaluating closed outgoing transport endpoint_id={} connection_id={} record_stable_id={:?} record_generation={} current_live_stable_id={:?} closed_transport_stable_id={:?} state={:?} status_reason={:?}",
endpoint_id_str,
record.connection_id,
record.transport_stable_id,
record.transport_generation,
current_live_stable_id,
closed_transport_stable_id,
record.state,
record.status_reason
)));
let replacement_window_open = matches!(
record.state,
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting
| crate::connection_manager::ConnectionState::Connected
);
if replacement_window_open {
let live_replacement_stable_id =
current_live_stable_id.filter(|stable_id| {
Some(*stable_id) != closed_transport_stable_id
});
if let Some(live_replacement_stable_id) = live_replacement_stable_id
{
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] closed outgoing transport rebound onto live replacement endpoint_id={} connection_id={} closed_transport_stable_id={:?} rebound_stable_id={:?} error={:?} active_state={:?}",
endpoint_id_str,
record.connection_id,
closed_transport_stable_id,
live_replacement_stable_id,
error,
record.state
)));
self.connection_manager
.mark_transport_replaced_by_if_current(
&record.connection_id,
closed_transport_stable_id
.expect("records are filtered by a stable id"),
live_replacement_stable_id,
Some("wasm-connect-bridge-closed-rebind".to_string()),
)
.await;
let _ = self
.confirm_managed_connection_readiness(&record.connection_id)
.await;
} else if terminal_disconnect_close {
let reason = error.clone().or_else(|| {
Some("wasm-connect-bridge-terminal-close".to_string())
});
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] terminal outgoing close settled closed state endpoint_id={} connection_id={} closed_transport_stable_id={:?} error={:?} active_state={:?}",
endpoint_id_str,
record.connection_id,
closed_transport_stable_id,
error,
record.state
)));
let closed = self
.connection_manager
.set_closed_if_current(
&record.connection_id,
closed_transport_stable_id
.expect("records are filtered by a stable id"),
reason,
)
.await;
if closed.is_some()
&& is_manual_disconnect_close_reason(error.as_deref())
{
self.suppress_auto_connect_for_connection_peer(
&record.connection_id,
"remote manual disconnect",
)
.await;
}
} else {
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] closed outgoing transport enters replacement-pending state endpoint_id={} connection_id={} closed_transport_stable_id={:?} error={:?} active_state={:?}",
endpoint_id_str,
record.connection_id,
closed_transport_stable_id,
error,
record.state
)));
self.connection_manager
.mark_transport_replaced_if_current(
&record.connection_id,
closed_transport_stable_id
.expect("records are filtered by a stable id"),
Some("wasm-connect-bridge-closed".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS
.to_string(),
),
)
.await;
self.wake_browser_auto_connect();
}
self.emit_current_wasm_connection_state(&record.connection_id)
.await;
continue;
}
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] retiring endpoint records after closed outgoing transport endpoint_id={} records={} error={:?}",
endpoint_id_str,
retire_count,
error
)));
let _ = self
.connection_manager
.set_closed_if_current(
&record.connection_id,
closed_transport_stable_id
.expect("records are filtered by a stable id"),
error
.clone()
.or_else(|| Some("wasm-connect-bridge-closed".to_string())),
)
.await;
}
}
}
}
});
}
pub async fn is_connected(&self, endpoint_id: iroh::EndpointId) -> bool {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
if node.is_connected(endpoint_id).await {
return true;
}
}
false
}
pub async fn is_connected_str(&self, endpoint_id: &str) -> anyhow::Result<bool> {
let endpoint_id = endpoint_id
.trim()
.parse::<iroh::EndpointId>()
.map_err(|error| anyhow::anyhow!("invalid peer nodeId: {error}"))?;
Ok(self.is_connected(endpoint_id).await)
}
}
impl Client {
#[cfg(not(target_arch = "wasm32"))]
pub fn spawn_iroh_path_watcher(
&self,
connection_id: &str,
remote_node_id: &str,
connection: &iroh::endpoint::Connection,
force_restart: bool,
) {
let stable_id = crate::transport_generation::for_connection(&connection);
if !self
.iroh_path_watcher_stable_ids
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.insert(stable_id)
{
return;
}
let client = self.clone();
let connection_id = connection_id.to_string();
let remote_node_id = remote_node_id.to_string();
let connection = connection.clone();
tokio::spawn(async move {
use futures::StreamExt as _;
let mut path_events = connection.path_events();
let initial_path_kind = {
let kind = client
.classify_native_iroh_connection_path(&connection)
.await;
if kind == IrohPathKind::Unknown {
None
} else {
Some(kind)
}
};
if let Some(path_kind) = initial_path_kind {
handle_path_change(
&client,
&connection_id,
&remote_node_id,
stable_id,
path_kind,
force_restart,
)
.await;
}
loop {
tokio::select! {
path_event = path_events.next() => {
if path_event.is_none() {
break;
}
let paths = connection.paths();
if !paths.iter().any(|p| p.is_selected()) {
continue;
}
let path_kind = client
.classify_native_iroh_connection_path(&connection)
.await;
if path_kind == IrohPathKind::Unknown {
continue;
}
handle_path_change(
&client,
&connection_id,
&remote_node_id,
stable_id,
path_kind,
false,
)
.await;
}
_ = connection.closed() => {
let close_reason = connection.close_reason();
let replacement_observed = client
.observe_native_transport_generation_lost(
&connection_id,
stable_id,
"iroh-connection-closed",
)
.await;
let close_reason_debug = format!("{:?}", close_reason);
if replacement_observed
&& !is_terminal_close_reason(Some(&close_reason_debug))
{
let _ = client
.schedule_native_external_auto_connect_recovery(
&remote_node_id,
)
.await;
}
break;
}
}
}
client
.iroh_path_watcher_stable_ids
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.remove(&stable_id);
});
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(super) async fn report_iroh_path_observation_if_current(
client: &super::Client,
connection_id: &str,
transport_stable_id: u64,
transport_generation: u64,
route_generation: u64,
active_transport: &str,
parallel_transport: Option<&str>,
) -> bool {
client
.report_transport_status_for_generation(
connection_id,
active_transport,
parallel_transport,
transport_stable_id,
transport_generation,
route_generation,
)
.await
.is_some()
}
#[cfg(not(target_arch = "wasm32"))]
pub(super) async fn handle_path_change(
client: &super::Client,
connection_id: &str,
remote_node_id: &str,
transport_stable_id: u64,
path_kind: IrohPathKind,
force_restart: bool,
) {
let Some(observation_record) = client
.connection_manager
.get_by_connection_id(connection_id)
.await
else {
return;
};
if observation_record.transport_stable_id != Some(transport_stable_id) {
eprintln!(
"[IrohPath] stale path observation ignored connection_id={} remote_node_id={} transport_stable_id={} path={:?}",
connection_id, remote_node_id, transport_stable_id, path_kind,
);
return;
}
let observation_transport_generation = observation_record.transport_generation;
let observation_route_generation = observation_record.route_generation;
if path_kind.is_relay_path() {
let remote_capabilities = client
.native_peer_transport_capabilities
.read()
.await
.get(connection_id)
.cloned()
.unwrap_or_default();
eprintln!(
"[IrohPath] relay path detected connection_id={} remote_node_id={} force_restart={} → triggering transport upgrade",
connection_id, remote_node_id, force_restart
);
if let Err(error) = client
.maybe_start_preferred_native_iroh_carrier(
connection_id,
Some(remote_node_id),
remote_capabilities.contains(&NativePeerTransportCapability::WebRtc),
remote_capabilities.contains(&NativePeerTransportCapability::Moq),
)
.await
{
eprintln!(
"[IrohPath] preferred carrier trigger failed connection_id={} error={}",
connection_id, error
);
}
let remote_supports_ble = remote_capabilities.contains(&NativePeerTransportCapability::Ble);
if remote_supports_ble && client.is_ble_transport_enabled().await {
if let Err(error) = client
.maybe_start_native_ble_upgrade(connection_id, Some(remote_node_id))
.await
{
eprintln!(
"[IrohPath] BLE upgrade trigger failed connection_id={} error={}",
connection_id, error
);
}
}
let _ = report_iroh_path_observation_if_current(
client,
connection_id,
transport_stable_id,
observation_transport_generation,
observation_route_generation,
"iroh-relay",
None,
)
.await;
} else {
eprintln!(
"[IrohPath] direct path detected kind={:?} connection_id={} remote_node_id={}",
path_kind, connection_id, remote_node_id
);
let _ = report_iroh_path_observation_if_current(
client,
connection_id,
transport_stable_id,
observation_transport_generation,
observation_route_generation,
path_kind.transport_label(),
None,
)
.await;
}
}