use super::*;
#[cfg(not(target_arch = "wasm32"))]
use iroh::Watcher;
#[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::reason_is_replacement_churn(error)
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn is_graceful_disconnect_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::reason_is_graceful_disconnect(error)
}
#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
pub(crate) fn is_terminal_disconnect_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::reason_is_terminal_disconnect(error)
}
pub(crate) fn is_manual_disconnect_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::reason_is_manual_disconnect(error)
}
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)
}
pub(crate) fn incoming_transport_requires_optional_route_restart(
already_managed_connected: bool,
replaced_native_main_route: bool,
) -> bool {
!already_managed_connected || replaced_native_main_route
}
#[cfg_attr(not(any(test, target_arch = "wasm32")), allow(dead_code))]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum WasmClosedTransportRebind {
TerminalDisconnect,
BaseReplacement(u64),
IndependentTransportOnly,
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>,
has_live_transport: bool,
) -> WasmClosedTransportRebind {
if terminal_disconnect {
return WasmClosedTransportRebind::TerminalDisconnect;
}
match current_base_stable_id {
Some(stable_id) => WasmClosedTransportRebind::BaseReplacement(stable_id),
None if has_live_transport => WasmClosedTransportRebind::IndependentTransportOnly,
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_compound_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"))]
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| connection.stable_id() as u64);
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 {
println!(
"[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))]
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(
project_id: String,
api_key: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> Self {
ClientBuilder::new(project_id, api_key, token_provider).build()
}
pub fn builder(
project_id: String,
api_key: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> ClientBuilder {
ClientBuilder::new(project_id, api_key, token_provider)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn builder_with_native_auth(
project_id: String,
api_key: String,
auth_state: crate::native_auth::NativeAuthState,
) -> ClientBuilder {
ClientBuilder::new(project_id, api_key, auth_state.token_provider())
.native_auth_state(auth_state)
}
pub fn builder_with_app_tag(
project_id: String,
app_tag: String,
token_provider: Box<dyn Fn() -> Option<String> + Send + Sync>,
) -> ClientBuilder {
ClientBuilder::new_with_app_tag(project_id, app_tag, token_provider)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn builder_with_native_space_auth(
project_id: String,
api_key: String,
space_key: String,
auth_state: crate::native_auth::NativeAuthState,
) -> ClientBuilder {
ClientBuilder::new_with_app_tag(
project_id,
crate::space_app_tag_from_keys(&api_key, &space_key),
auth_state.token_provider(),
)
.native_auth_state(auth_state)
.auth_mode(AuthMode::Anonymous)
}
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::ScopedConnectionActorRegistry>,
) {
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::ScopedConnectionActorRegistry>> {
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::NativeSystemDeviceInfo> {
let identity = self.get_native_device_identity().await?;
Ok(identity
.system_info
.unwrap_or_else(crate::native_device::current_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::NativeDeviceNameSource::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 subscribe_native_connection_state_updates(
&self,
) -> tokio::sync::broadcast::Receiver<crate::client::ConnectionStateSnapshot> {
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,
) -> 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 mut builder = iroh::Endpoint::builder(iroh::endpoint::presets::N0);
let mut alpns = extra_alpns;
alpns.push(b"plutonium/p2p/1".to_vec());
builder = builder.alpns(alpns);
let transport_config = self.transport_config.read().await.clone();
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::IrohRelayTransportPolicy::WebsocketRequired);
builder = apply_native_network_preferences(builder, relay_only, relay_transport_policy)?;
builder = self.apply_optional_lan_discovery(builder).await?;
builder = builder.secret_key(endpoint_secret_key);
let endpoint = builder.bind().await?;
let node_id = endpoint.id().to_string();
self.start_local_discovery_tasks(&endpoint, &node_id).await;
let relay_was_online_before_publish =
match tokio::time::timeout(std::time::Duration::from_secs(10), endpoint.online()).await
{
Ok(()) => {
println!("[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);
}
});
}
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 = spawn_internal_router.then(|| 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 let Some(events) = native_accept_events {
self.start_native_accept_bridge(events);
}
self.start_native_incoming_stream_router(native_incoming_streams);
#[cfg(not(target_arch = "wasm32"))]
println!(
"[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)
.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 mut alpns = extra_alpns;
alpns.push(b"plutonium/p2p/1".to_vec());
builder = builder.alpns(alpns).relay_mode(iroh::RelayMode::Default);
if let Some(key_bytes) = secret_key {
wasm_init_log("init_iroh:using-provided-secret-key");
let key = iroh::SecretKey::try_from(&key_bytes[..])?;
builder = builder.secret_key(key);
}
wasm_init_log("init_iroh:bind-start");
let endpoint = builder.bind().await?;
wasm_init_log("init_iroh:bind-complete");
let node_id = endpoint.id().to_string();
wasm_init_log("init_iroh:spawn-router-start");
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)
.await
}
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"))
}
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 = spawn_internal_router.then(|| node.accept_events());
let native_incoming_streams = node.incoming_streams_stream();
*self.iroh_node.write().await = Some(node);
if let Some(events) = native_accept_events {
self.start_native_accept_bridge(events);
}
self.start_native_incoming_stream_router(native_incoming_streams);
println!(
"[pluto-rtc][native] adopted endpoint node_id={} internal_router={}",
node_id, spawn_internal_router
);
Ok(node_id)
}
#[cfg(not(target_arch = "wasm32"))]
pub 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,
_ => 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 register_native_transport_upgrade_provider(
&self,
provider: Arc<dyn crate::client::NativeTransportUpgradeProvider>,
) -> 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::NativeTransportUpgradeProvider>> {
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_compound_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());
println!(
"[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_compound_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());
println!(
"[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());
println!(
"[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_compound_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());
println!(
"[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?;
println!(
"[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_compound_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());
println!(
"[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::build_compound_ticket_with_expiry(
&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());
println!(
"[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 mark_endpoint_dial_timed_out(
&self,
endpoint_id: &iroh::EndpointId,
label: &'static str,
) {
use crate::connection_manager::ConnectionState;
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;
}
let reason = format!("{label}-timeout");
for record in records {
if matches!(
record.state,
ConnectionState::Pending | ConnectionState::Connecting
) {
let _ = self
.connection_manager
.set_failed(&record.connection_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 (send, recv) =
self.wrap_peer_streams_for_connection(connection_id.as_deref(), 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());
}
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| connection.stable_id() as u64)
else {
return Ok(connection_id);
};
self.commit_current_native_transport_record(
&connection_id,
&remote_node_id,
device_id_hint.clone(),
transport_stable_id,
source.clone(),
)
.await;
let committed_transport_is_current = self
.get_connection(endpoint_id)
.await
.is_some_and(|connection| connection.stable_id() as u64 == transport_stable_id);
if committed_transport_is_current {
return Ok(connection_id);
}
println!(
"[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: u64,
transport_source: Option<String>,
) {
#[cfg(not(target_arch = "wasm32"))]
self.retire_stale_native_main_route(connection_id, transport_stable_id)
.await;
if self
.connection_manager
.get_by_connection_id(connection_id)
.await
.is_none()
{
self.connection_manager
.upsert_pending(
connection_id.to_string(),
Some(remote_node_id.to_string()),
device_id_hint,
Some(remote_node_id.to_string()),
)
.await;
}
self.connection_manager
.set_connected_with_transport(
connection_id,
Some(remote_node_id.to_string()),
Some(transport_stable_id),
transport_source,
)
.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;
}
if self.trusted_user_device_application_crypto_is_required()
&& !self.connection_application_crypto_is_confirmed(connection_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 needs_commit = match self
.connection_manager
.get_by_connection_id(&connection_id)
.await
{
None => true,
Some(record) => !matches!(record.state, ConnectionState::Connected),
};
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 {
println!(
"[pluto-rtc][iroh] redial using cached endpoint addr endpoint_id={}",
endpoint_id
);
return self.ensure_connected_addr(endpoint_id, endpoint_addr).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(&endpoint_id, "ensure_connected")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next() => {
match event {
Some(ConnectEvent::Connected) => {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected".to_string()),
)
.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;
}
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(&endpoint_id, "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.emit_current_wasm_connection_state(&connection_id).await;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64);
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 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(&endpoint_id, "ensure_connected_addr")
.await;
return Err(anyhow::anyhow!("Timed out connecting to {}", endpoint_id));
}
event = events.next() => {
match event {
Some(ConnectEvent::Connected) => {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr".to_string()),
)
.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"))]
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) = match event {
AcceptEvent::Accepted {
endpoint_id,
transport_stable_id,
} => (endpoint_id, transport_stable_id),
AcceptEvent::Closed {
endpoint_id,
transport_stable_id,
error,
..
} => {
let endpoint_id_str = endpoint_id.to_string();
let records = client
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
for record in records {
if client
.observe_native_transport_generation_lost(
&record.connection_id,
transport_stable_id,
"native-accept-closed",
)
.await
{
println!(
"[pluto-rtc][native-accept] current transport closed connection_id={} endpoint_id={} transport_stable_id={} error={:?}",
record.connection_id,
endpoint_id,
transport_stable_id,
error,
);
client
.emit_current_native_connection_state(&record.connection_id)
.await;
}
}
continue;
}
};
let current_stable_id = client
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64);
if !accepted_transport_event_is_current(transport_stable_id, current_stable_id) {
println!(
"[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;
}
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
{
println!(
"[pluto-rtc][native-accept] promoted trusted user-device connection_id={} remote_node_id={} device_id={}",
connection_id, remote_node_id, device_id
);
}
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>,
) {
let client = self.clone();
let application_streams = self.native_application_streams.clone();
tokio::spawn(async move {
while let Ok(incoming_stream) = incoming.recv().await {
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(send, recv) => {
match client
.route_incoming_bi_stream_for_admission(
endpoint_id,
transport_stable_id,
send,
recv,
)
.await
{
Ok(crate::client::IncomingBiStreamDisposition::Consumed) => {}
Ok(crate::client::IncomingBiStreamDisposition::Forward {
send,
recv,
}) => {
if application_streams
.send(IncomingStream {
endpoint_id,
transport_stable_id,
stream: crate::native_node::IncomingStreamType::Bi(
send, recv,
),
})
.await
.is_err()
{
break;
}
}
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,
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 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(&endpoint_id, "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.emit_current_wasm_connection_state(&connection_id).await;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64);
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.get_connection(endpoint_id)
.await
.is_some_and(|connection| connection.stable_id() as u64 == expected_transport_stable_id)
}
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() {
node.disconnect_with_reason(endpoint_id, reason).await?;
let records = self
.connection_manager
.get_by_endpoint_id(&endpoint_id_str)
.await;
for record in records {
if crate::lifecycle_reason::reason_is_transient_reconnect(Some(reason)) {
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;
} else {
self.retire_managed_connection(&record.connection_id, Some(reason.to_string()))
.await;
}
}
Ok(())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[cfg(not(target_arch = "wasm32"))]
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);
}
for record in records {
self.observe_native_transport_generation_lost(
&record.connection_id,
expected_transport_stable_id,
"liveness-probe-stale",
)
.await;
}
Ok(true)
}
#[cfg(not(target_arch = "wasm32"))]
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
}
pub fn runtime_policy_snapshot(&self) -> crate::runtime_policy::RuntimePolicySnapshot {
crate::runtime_policy::runtime_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
)
})
}
}
}
}
#[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) => {
println!(
"[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
)
})
}
}
}
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"))]
println!(
"[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() {
println!(
"[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"))]
println!("{}", 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());
}
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_republish_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();
println!(
"[pluto-rtc][auto-connect][presence-republish] user_id={} local_device_id={} reason={} owner=native-presence-actor queued={} liveness_source=rtdb",
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) => {
println!(
"[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_republish_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 = connection.stable_id() as u64;
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 mut accept_connection =
Box::pin(node.accept_external_connection_with_install_notifier(
connection.clone(),
install_outcome_sender,
));
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 => {
result?;
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.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.await;
println!(
"[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;
}
}
let already_managed_connected = self
.connection_manager
.peer_snapshot(&connection_id)
.await
.map(|s| {
matches!(
s.status,
crate::connection_manager::ConnectionState::Connected
)
})
.unwrap_or(false);
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(|active| active.stable_id() as u64);
if active_stable_id != Some(stable_id) {
println!(
"[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.await;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
return result;
}
let replaced_native_main_route =
!self.native_admission_route_is_ready_for_transport(&connection_id, Some(stable_id));
let _ = self
.promote_known_native_user_device_connection(&connection_id, &remote_node_id)
.await;
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
let force_optional_route_restart = incoming_transport_requires_optional_route_restart(
already_managed_connected,
replaced_native_main_route,
);
if force_optional_route_restart {
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
if let Err(error) = self
.request_native_webrtc_recovery(
&connection_id,
Some(&remote_node_id),
crate::native_webrtc_policy::NativeWebRTCRecoveryTrigger::Native(
crate::native_webrtc_policy::NativeWebRTCNativeTrigger::IncomingConnection,
),
crate::native_webrtc_policy::NativeWebRTCRecoveryOptions {
force_restart: true,
preferred_negotiation_id: None,
role_override: None,
},
)
.await
{
eprintln!(
"[NativeWebRTC] fallback trigger failed source=incoming-connection connection_id={} remote_node_id={} error={}",
connection_id,
remote_node_id,
error
);
}
}
}
println!(
"[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 { .. }
) {
println!(
"[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);
}
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 = timeout_connection.stable_id() as u64;
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 {
println!(
"[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;
}
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
{
if client
.is_managed_retirement_deferred(&connection_id_for_timeout)
.await
{
return;
}
let webrtc_active = matches!(
client
.native_webrtc_state_for_peer(&connection_id_for_timeout)
.await,
Some((_, crate::transport::NativeWebRTCState::Connecting))
| Some((_, crate::transport::NativeWebRTCState::Connected))
);
if webrtc_active {
if !client
.admission_timeout_still_owns_transport(
&connection_id_for_timeout,
remote_endpoint_id_for_timeout,
timeout_transport_stable_id,
)
.await
{
return;
}
client
.retire_managed_connection(
&connection_id_for_timeout,
Some("session-admission-timeout".to_string()),
)
.await;
println!(
"[PlutoRTC] handle_incoming_connection admission timeout deferred: WebRTC still active connection_id={} remote_node_id={}",
connection_id_for_timeout,
remote_node_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;
}
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;
println!(
"[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_optional_route_restart,
);
let result = accept_connection.await;
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);
println!(
"[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(
&connection_id,
stable_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
if closed.is_none() {
println!(
"[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 => {
println!(
"[PlutoRTC] handle_incoming_connection: {} deduplicated (locally closed); \
preserving active manager record for kept connection",
connection_id,
);
}
IncomingTransportCloseResolution::RetiredClosedTransport => {
println!(
"[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 !self
.connection_manager
.current_transport_matches(connection_id, Some(closed_stable_id))
.await
{
if let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
{
println!(
"[PlutoRTC] handle_incoming_connection ignoring stale closed transport connection_id={} remote_node_id={} closed_stable_id={} active_stable_id={:?} active_generation={} status_reason={:?}",
connection_id,
remote_node_id,
closed_stable_id,
record.transport_stable_id,
record.transport_generation,
record.status_reason,
);
}
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
if is_manual_disconnect_close_reason(Some(close_reason_debug)) {
let closed = self
.connection_manager
.set_closed_if_current(
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 = current.stable_id() as u64;
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 {
println!(
"[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;
}
println!(
"[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 incoming_local_close_should_wait_for_replacement(close_reason_debug) {
println!(
"[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 let Some(record) = self
.connection_manager
.get_by_connection_id(connection_id)
.await
{
let replacement_pending = matches!(
record.state,
crate::connection_manager::ConnectionState::Pending
| crate::connection_manager::ConnectionState::Connecting
| crate::connection_manager::ConnectionState::Connected
) && matches!(
record.status_reason.as_deref(),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS)
| Some("incoming-kept-existing")
);
if replacement_pending {
println!(
"[PlutoRTC] handle_incoming_connection preserving replacement-pending manager record connection_id={} remote_node_id={} closed_stable_id={} active_stable_id={:?} active_generation={} status_reason={:?}",
connection_id,
remote_node_id,
closed_stable_id,
record.transport_stable_id,
record.transport_generation,
record.status_reason,
);
return IncomingTransportCloseResolution::PreservedKeptTransport;
}
}
println!(
"[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 closed = self
.connection_manager
.set_closed_if_current(
connection_id,
closed_stable_id,
Some(
crate::lifecycle_reason::REASON_INCOMING_TRANSPORT_CLOSED_WITHOUT_LIVE_REPLACEMENT
.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"))]
println!(
"[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::build_framed_manual_disconnect();
let Ok(mut send) = self.open_uni(endpoint_id).await else {
return;
};
if tokio::io::AsyncWriteExt::write_all(&mut send, &frame)
.await
.is_ok()
{
let _ = send.finish();
}
}
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
}
}
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,
};
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(),
);
}
#[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 {
IrohPathKind::Relay
}
#[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| connection.stable_id() as u64)
.ok_or_else(|| {
anyhow::anyhow!(
"cannot promote WASM base transport without a physical stable id for {}",
remote_node_id
)
})?;
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);
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
None,
Some(remote_node_id.clone()),
)
.await;
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.clone()),
Some(stable_id),
Some(source.to_string()),
)
.await;
let _ = self
.report_transport_status_for_current_generation(&connection_id, "iroh-relay", 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| connection.stable_id() as u64);
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 Err(error) = self
.promote_wasm_incoming_transport(endpoint_id, "wasm-accept")
.await
{
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] failed promoting accepted transport endpoint_id={} error={}",
endpoint_id, error
)));
}
}
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 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 is_terminal_disconnect_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 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;
}
let current_live_stable_id = self
.get_connection(endpoint_id)
.await
.map(|conn| conn.stable_id() as u64);
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] accepted transport awaiting replacement after local close endpoint_id={} connection_id={} current_live_stable_id={:?} error={:?}",
remote_node_id,
connection_id,
current_live_stable_id,
error
)));
self.connection_manager
.mark_transport_replaced_if_current(
&connection_id,
transport_stable_id,
Some("wasm-accept-closed".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string(),
),
)
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
} else {
let new_stable_id = self
.get_connection(endpoint_id)
.await
.map(|conn| conn.stable_id() as u64);
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| connection.stable_id() as u64);
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_disconnect_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| connection.stable_id() as u64);
match classify_wasm_closed_transport_rebind(
terminal_disconnect_close,
current_live_stable_id,
live_transport,
) {
WasmClosedTransportRebind::TerminalDisconnect => {
}
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::IndependentTransportOnly => {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] base transport closed while independent route remains; awaiting concrete base 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.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;
}
}
let endpoint_id = endpoint_id.to_string();
self.connection_manager
.get_by_endpoint_id(&endpoint_id)
.await
.into_iter()
.any(|record| {
matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) && crate::transport_label::is_independent_transport(
record.active_transport.as_str(),
)
})
}
#[cfg(all(
not(target_arch = "wasm32"),
not(any(target_os = "ios", target_os = "android"))
))]
pub async fn fetch_and_cache_app_limits(&self, api_key: &str) -> anyhow::Result<()> {
use crate::firebase::schema::developer_app_doc;
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use firestore::{FirestoreDb, FirestoreDbOptions};
use gcloud_sdk::{ExternalJwtFunctionSource, Token, TokenSourceType};
fn parse_id_token_expiry_utc(id_token: &str) -> Option<chrono::DateTime<chrono::Utc>> {
let mut parts = id_token.split('.');
let _header = parts.next()?;
let payload = parts.next()?;
let payload_bytes = URL_SAFE_NO_PAD.decode(payload).ok()?;
let payload_json: serde_json::Value = serde_json::from_slice(&payload_bytes).ok()?;
let exp_seconds = payload_json.get("exp")?.as_i64()?;
chrono::DateTime::<chrono::Utc>::from_timestamp(exp_seconds, 0)
}
let db = {
let tp = self.token_provider.clone();
if tp().is_some_and(|t| !t.trim().is_empty()) {
let tp2 = tp.clone();
let token_source = ExternalJwtFunctionSource::new(move || {
let tp3 = tp2.clone();
async move {
let token = tp3()
.filter(|t| !t.trim().is_empty())
.ok_or_else(|| gcloud_sdk::error::ErrorKind::TokenSource)?;
let expires_at = parse_id_token_expiry_utc(&token)
.unwrap_or_else(|| chrono::Utc::now() + chrono::Duration::minutes(10));
Ok(Token::new("Bearer".to_string(), token.into(), expires_at))
}
});
let options = FirestoreDbOptions::new(self.project_id.clone());
FirestoreDb::with_options_token_source(
options,
gcloud_sdk::GCP_DEFAULT_SCOPES.clone(),
TokenSourceType::ExternalSource(Box::new(token_source)),
)
.await?
} else {
FirestoreDb::new(&self.project_id).await?
}
};
let path = developer_app_doc(api_key);
let doc: Option<serde_json::Value> = db
.fluent()
.select()
.by_id_in("developer_apps")
.obj()
.one(api_key)
.await?;
let Some(doc) = doc else {
anyhow::bail!("developer_apps/{} not found", api_key);
};
fn read_i64(val: &serde_json::Value, key: &str) -> i64 {
val.get(key)
.and_then(|v| {
v.as_i64()
.or_else(|| v.as_str().and_then(|s| s.parse().ok()))
})
.unwrap_or(-1)
}
let limits_val = doc
.get("limits")
.cloned()
.unwrap_or(serde_json::Value::Null);
let limits = crate::client::AppLimits {
devices_per_user: read_i64(&limits_val, "devicesPerUser"),
max_rooms: read_i64(&limits_val, "maxRooms"),
max_members_per_room: read_i64(&limits_val, "maxMembersPerRoom"),
};
*self.app_limits.write().await = limits;
let _ = path; Ok(())
}
#[cfg(any(target_os = "ios", target_os = "android"))]
pub async fn fetch_and_cache_app_limits(&self, _api_key: &str) -> anyhow::Result<()> {
Ok(())
}
}
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 = connection.stable_id() as u64;
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 = connection.paths();
let initial_path_kind = if initial.iter().any(|p| p.is_selected()) {
Some(client.iroh_path_kind(&remote_node_id).await)
} else {
let kind = client.iroh_path_kind(&remote_node_id).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.iroh_path_kind(&remote_node_id).await;
if path_kind == IrohPathKind::Unknown {
continue;
}
handle_path_change(
&client,
&connection_id,
&remote_node_id,
stable_id,
path_kind,
false,
)
.await;
}
_ = connection.closed() => {
client
.observe_native_transport_generation_lost(
&connection_id,
stable_id,
"iroh-connection-closed",
)
.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) {
println!(
"[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() {
#[cfg(feature = "transport-webrtc")]
let _ = client
.clear_native_webrtc_suppression(connection_id, Some("iroh-quic-primary"))
.await;
println!(
"[IrohPath] relay path detected connection_id={} remote_node_id={} force_restart={} → triggering transport upgrade",
connection_id, remote_node_id, force_restart
);
if client.is_webrtc_transport_enabled().await {
#[cfg(all(not(target_arch = "wasm32"), feature = "transport-webrtc"))]
if let Err(e) = client
.request_native_webrtc_recovery(
connection_id,
Some(remote_node_id),
crate::native_webrtc_policy::NativeWebRTCRecoveryTrigger::Native(
crate::native_webrtc_policy::NativeWebRTCNativeTrigger::RelayPath,
),
crate::native_webrtc_policy::NativeWebRTCRecoveryOptions {
force_restart,
preferred_negotiation_id: None,
role_override: None,
},
)
.await
{
println!(
"[IrohPath] webrtc upgrade trigger failed connection_id={} error={}",
connection_id, e
);
}
}
if client.is_moq_transport_enabled().await {
if let Err(e) = client
.maybe_start_native_moq_upgrade(connection_id, Some(remote_node_id))
.await
{
println!(
"[IrohPath] moq upgrade trigger failed connection_id={} error={}",
connection_id, e
);
}
}
let remote_supports_ble = client
.native_peer_transport_capabilities
.read()
.await
.get(connection_id)
.is_some_and(|capabilities| capabilities.contains(&IrohPathKind::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
{
println!(
"[IrohPath] BLE upgrade trigger failed connection_id={} error={}",
connection_id, error
);
}
}
#[cfg(feature = "transport-webrtc")]
{
let connected_webrtc_session = if let Some(session) = client
.get_connected_webrtc_session_for_peer(connection_id)
.await
{
Some(session)
} else {
client
.get_connected_webrtc_session_for_peer(remote_node_id)
.await
};
if let Some(session) = connected_webrtc_session {
let webrtc_transport = session
.selected_ice_pair_summary()
.await
.map(|pair| pair.transport_label())
.unwrap_or(crate::transport_label::WEBRTC);
let (active_transport, parallel_transport) = if crate::transport_label::is_webrtc(
observation_record.active_transport.as_str(),
) {
(webrtc_transport, Some("iroh-relay"))
} else {
("iroh-relay", Some(webrtc_transport))
};
let _ = report_iroh_path_observation_if_current(
client,
connection_id,
transport_stable_id,
observation_transport_generation,
observation_route_generation,
active_transport,
parallel_transport,
)
.await;
} else {
let _ = report_iroh_path_observation_if_current(
client,
connection_id,
transport_stable_id,
observation_transport_generation,
observation_route_generation,
"iroh-relay",
None,
)
.await;
}
}
#[cfg(not(feature = "transport-webrtc"))]
{
let _ = report_iroh_path_observation_if_current(
client,
connection_id,
transport_stable_id,
observation_transport_generation,
observation_route_generation,
"iroh-relay",
None,
)
.await;
}
} else {
#[cfg(feature = "transport-webrtc")]
{
client
.suppress_native_webrtc_restarts(
connection_id,
super::Client::NATIVE_WEBRTC_DIRECT_QUIC_SUPPRESSION_MS,
path_kind.transport_label(),
)
.await;
let _ = client
.close_connecting_native_webrtc_session(connection_id, path_kind.transport_label())
.await;
}
println!(
"[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;
}
}