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)
}
pub(crate) fn is_manual_disconnect_close_reason(error: Option<&str>) -> bool {
crate::lifecycle_reason::reason_is_manual_disconnect(error)
}
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_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)
}
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)
}
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 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;
let node_guard = self.node_id.read().await;
if endpoint_guard.is_some() {
if let Some(existing_node_id) = node_guard.clone() {
return Ok(existing_node_id);
}
}
}
let endpoint_secret_key = match secret_key {
Some(key_bytes) => iroh::SecretKey::try_from(&key_bytes[..])?,
None => iroh::SecretKey::generate(),
};
let endpoint_id = endpoint_secret_key.public();
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::Auto);
builder = apply_native_network_preferences(builder, relay_only, relay_transport_policy)?;
builder = self.apply_optional_lan_discovery(builder).await?;
builder = self
.apply_optional_ble_transport(builder, endpoint_id)
.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;
match tokio::time::timeout(std::time::Duration::from_secs(10), endpoint.online()).await {
Ok(()) => {
println!("[pluto-rtc][native] relay online node_id={}", node_id);
}
Err(_) => {
eprintln!(
"[pluto-rtc][native] WARNING: relay not connected after 10s, \
tickets may lack relay URLs node_id={}",
node_id
);
}
}
#[cfg(not(target_arch = "wasm32"))]
{
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::UpdateNow);
}
});
}
let node = if spawn_internal_router {
IrohNativeNode::spawn_with_endpoint(endpoint.clone()).await?
} else {
IrohNativeNode::spawn_with_endpoint_no_router(endpoint.clone()).await?
};
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);
#[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
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn init_iroh_ble_hardware_harness(
&self,
secret_key: Option<Vec<u8>>,
extra_alpns: Vec<Vec<u8>>,
) -> anyhow::Result<String> {
let _init_lock = self.iroh_init_guard.lock().await;
{
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() {
return Ok(existing_node_id);
}
}
}
let endpoint_secret_key = match secret_key {
Some(key_bytes) => iroh::SecretKey::try_from(&key_bytes[..])?,
None => iroh::SecretKey::generate(),
};
let endpoint_id = endpoint_secret_key.public();
let ble = Self::build_ble_transport(
endpoint_id,
crate::client::BleConfig {
enabled: true,
connect_timeout_ms: Some(20_000),
retry_attempts: Some(15),
retry_backoff_ms: Some(500),
},
)
.await?;
*self.ble_transport.write().await = Some(ble.clone());
let mut alpns = extra_alpns;
alpns.push(b"plutonium/p2p/1".to_vec());
let idle_timeout: iroh::endpoint::IdleTimeout = std::time::Duration::from_secs(15)
.try_into()
.map_err(|error| anyhow::anyhow!("invalid BLE idle timeout: {error}"))?;
let transport_config = iroh::endpoint::QuicTransportConfig::builder()
.max_idle_timeout(Some(idle_timeout))
.build();
let endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::N0DisableRelay)
.alpns(alpns)
.hooks(ble.dedup_hook())
.add_custom_transport(ble.as_custom_transport())
.address_lookup(ble.address_lookup())
.transport_config(transport_config)
.secret_key(endpoint_secret_key)
.clear_ip_transports()
.bind()
.await?;
let node_id = endpoint.id().to_string();
let node = IrohNativeNode::spawn_with_endpoint(endpoint.clone()).await?;
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 endpoint_guard = self.iroh_endpoint.write().await;
*endpoint_guard = Some(endpoint);
println!(
"[pluto-rtc][ble-hardware-harness] initialized BLE-only endpoint node_id={}",
node_id
);
Ok(node_id)
}
pub async fn get_endpoint(&self) -> anyhow::Result<Endpoint> {
let guard = self.iroh_endpoint.read().await;
guard
.clone()
.ok_or_else(|| anyhow::anyhow!("Iroh Endpoint not initialized"))
}
pub async fn current_node_id(&self) -> Option<String> {
self.node_id.read().await.clone()
}
pub async fn adopt_endpoint(&self, endpoint: Endpoint) {
let node_id = endpoint.id().to_string();
{
let endpoint_guard = self.iroh_endpoint.read().await;
let node_guard = self.node_id.read().await;
let already_same_endpoint = endpoint_guard
.as_ref()
.map(|existing| existing.id().to_string() == node_id)
.unwrap_or(false)
&& node_guard
.as_ref()
.map(|existing| existing == &node_id)
.unwrap_or(false);
if already_same_endpoint {
return;
}
}
let mut node_guard = self.node_id.write().await;
*node_guard = Some(node_id);
drop(node_guard);
let mut endpoint_guard = self.iroh_endpoint.write().await;
*endpoint_guard = Some(endpoint);
#[cfg(not(target_arch = "wasm32"))]
{
let endpoint = endpoint_guard.clone();
drop(endpoint_guard);
if let Some(endpoint) = endpoint {
match IrohNativeNode::spawn_with_endpoint(endpoint).await {
Ok(node) => {
let mut node_guard = self.iroh_node.write().await;
*node_guard = Some(node);
}
Err(error) => {
eprintln!(
"[pluto-rtc][native] failed to adopt native iroh node runtime: {}",
error
);
}
}
}
}
}
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
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
async fn maybe_prime_ble_target_scan(&self, endpoint_id: iroh::EndpointId) {
let Some(transport) = self.ble_transport.read().await.clone() else {
return;
};
if transport.has_scan_hint_for_endpoint(&endpoint_id) {
match transport.ensure_connecting_for_endpoint(endpoint_id).await {
Ok(true) => {
println!(
"[pluto-rtc][ble] connection nudge sent endpoint_id={}",
endpoint_id
);
}
Ok(false) => {}
Err(error) => {
eprintln!(
"[pluto-rtc][ble] connection nudge failed endpoint_id={} error={}",
endpoint_id, error
);
}
}
return;
}
match transport.scan_for_endpoint(endpoint_id).await {
Ok(service_uuid) => {
println!(
"[pluto-rtc][ble] targeted scan primed endpoint_id={} service_uuid={}",
endpoint_id, service_uuid
);
}
Err(error) => {
eprintln!(
"[pluto-rtc][ble] targeted scan failed endpoint_id={} error={}",
endpoint_id, error
);
}
}
}
#[cfg(all(not(target_arch = "wasm32"), not(openrtc_unpublished_ble)))]
async fn maybe_prime_ble_target_scan(&self, _endpoint_id: iroh::EndpointId) {}
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;
let transport_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64);
if self
.connection_manager
.get_by_connection_id(&connection_id)
.await
.is_none()
{
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
device_id_hint.clone(),
Some(remote_node_id.clone()),
)
.await;
}
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id),
transport_stable_id,
transport_source,
)
.await;
Ok(connection_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.active_endpoint_ids().await.contains(&endpoint_id)
};
if already_active {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected-active".to_string()),
)
.await?;
return Ok(());
}
self.maybe_prime_ble_target_scan(endpoint_id).await;
let mut events = {
let node_guard = self.iroh_node.read().await;
let Some(node) = node_guard.as_ref() else {
return Err(anyhow::anyhow!("Iroh node not initialized"));
};
node.connect(endpoint_id)
};
let timeout = tokio::time::sleep(timeout);
tokio::pin!(timeout);
loop {
tokio::select! {
_ = &mut timeout => {
self.mark_endpoint_dial_timed_out(&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"));
};
let active = node.active_endpoint_ids().await;
if active.contains(&endpoint_id) {
drop(node_guard);
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected-active".to_string()),
)
.await?;
return Ok(());
}
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) => {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected".to_string()),
)
.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.active_endpoint_ids().await.contains(&endpoint_id)
};
if already_active {
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr-active".to_string()),
)
.await?;
return Ok(());
}
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(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"));
};
let active = node.active_endpoint_ids().await;
if active.contains(&endpoint_id) {
drop(node_guard);
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr-active".to_string()),
)
.await?;
return Ok(());
}
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
)));
self.finalize_transport_dial_record(
endpoint_id,
Some("ensure_connected_addr".to_string()),
)
.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 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 {
self.retire_managed_connection(&record.connection_id, Some(reason.to_string()))
.await;
}
Ok(())
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
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
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn open_bi_internal(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.open_bi(endpoint_id).await
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[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
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn open_uni_internal(
&self,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::SendStream> {
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.open_uni(endpoint_id).await
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[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?;
let node_guard = self.iroh_node.read().await;
if let Some(node) = node_guard.as_ref() {
node.open_uni(endpoint_id).await
} else {
Err(anyhow::anyhow!("Iroh node not initialized"))
}
}
#[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 let Some(node) = node_guard.as_ref() {
println!(
"[pluto-rtc][native] incoming stream receiver attached node_id={}",
self.current_node_id().await.as_deref().unwrap_or("unknown")
);
Ok(node.incoming_streams_stream())
} 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)
}
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn register_ble_hardware_harness_peer(
&self,
remote_node_id: &str,
application_crypto_key: [u8; crate::application_crypto::APPLICATION_KEY_BYTES],
) -> anyhow::Result<String> {
let local_node_id = self
.current_node_id()
.await
.ok_or_else(|| anyhow::anyhow!("Iroh node not initialized"))?;
let remote_node_id = remote_node_id.trim();
if remote_node_id.is_empty() {
anyhow::bail!("remote node id cannot be empty");
}
let endpoint_id = remote_node_id.parse::<iroh::EndpointId>()?;
let connection = self
.get_connection(endpoint_id)
.await
.ok_or_else(|| anyhow::anyhow!("no active BLE transport for {remote_node_id}"))?;
let stable_id = connection.stable_id() as u64;
let connection_id = Self::deterministic_connection_id(&local_node_id, remote_node_id);
self.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.to_string()),
Some(format!("ble-hardware-{remote_node_id}")),
Some(remote_node_id.to_string()),
)
.await;
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.to_string()),
Some(stable_id),
Some("ble-hardware-harness".to_string()),
)
.await
.ok_or_else(|| anyhow::anyhow!("failed to register connection {connection_id}"))?;
self.connection_manager
.set_health(
&connection_id,
crate::connection_manager::ConnectionHealth::Healthy,
)
.await;
self.set_connection_application_crypto_key(&connection_id, application_crypto_key);
self.set_connection_application_crypto_required(&connection_id);
Ok(connection_id)
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn ble_hardware_harness_debug(
&self,
target_node_id: Option<&str>,
) -> BleHardwareHarnessDebugSnapshot {
let enabled = self.ble_transport_enabled().await;
let local_node_id = self.current_node_id().await;
let target = target_node_id
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned);
let target_endpoint = target
.as_deref()
.and_then(|value| value.parse::<iroh::EndpointId>().ok());
let Some(transport) = self.ble_transport.read().await.clone() else {
return BleHardwareHarnessDebugSnapshot {
enabled,
local_node_id,
target_node_id: target,
local_service_uuid: None,
target_service_uuid: target_endpoint.as_ref().map(|endpoint| {
iroh_ble_transport::BleTransport::service_uuid_for_endpoint(endpoint)
.to_string()
}),
target_seen: false,
route_pipes: 0,
route_pipe_tombstones: 0,
route_scan_hints: 0,
route_pending: 0,
route_routable: 0,
route_reservations: 0,
peers: Vec::new(),
note: Some("BLE transport has not been initialized".to_string()),
};
};
let routes = transport.routing_snapshot();
let local_service_uuid = Some(transport.advertised_service_uuid().to_string());
let target_service_uuid = target_endpoint.as_ref().map(|endpoint| {
iroh_ble_transport::BleTransport::service_uuid_for_endpoint(endpoint).to_string()
});
let peers = transport
.snapshot_peers()
.into_iter()
.map(|peer| BleHardwareHarnessPeerSnapshot {
device_id: peer.device_id.to_string(),
phase: format!("{:?}", peer.phase),
phase_detail: peer.phase_detail,
consecutive_failures: peer.consecutive_failures,
connect_path: peer.connect_path.map(|path| format!("{path:?}")),
verified_endpoint: peer.verified_endpoint.map(|endpoint| endpoint.to_string()),
})
.collect::<Vec<_>>();
let target_seen = target_endpoint
.as_ref()
.map(|endpoint| transport.has_scan_hint_for_endpoint(endpoint))
.unwrap_or(false);
let note = match (&target, &target_endpoint) {
(Some(value), None) => Some(format!(
"target node id could not be parsed as an iroh EndpointId: {value}"
)),
_ => None,
};
BleHardwareHarnessDebugSnapshot {
enabled,
local_node_id,
target_node_id: target,
local_service_uuid,
target_service_uuid,
target_seen,
route_pipes: routes.pipes,
route_pipe_tombstones: routes.pipe_tombstones,
route_scan_hints: routes.scan_hints,
route_pending: routes.pending,
route_routable: routes.routable,
route_reservations: routes.reservations,
peers,
note,
}
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn pause_ble_hardware_harness_central_scan(&self) -> anyhow::Result<()> {
let transport = self
.ble_transport
.read()
.await
.clone()
.ok_or_else(|| anyhow::anyhow!("BLE transport has not been initialized"))?;
transport
.stop_central_scan()
.await
.map_err(|error| anyhow::anyhow!("failed to stop BLE central scan: {error}"))
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn prime_ble_hardware_harness_target_scan(
&self,
target_node_id: &str,
) -> anyhow::Result<String> {
let endpoint_id = target_node_id
.parse::<iroh::EndpointId>()
.map_err(|error| anyhow::anyhow!("invalid BLE target endpoint id: {error}"))?;
let transport = self
.ble_transport
.read()
.await
.clone()
.ok_or_else(|| anyhow::anyhow!("BLE transport has not been initialized"))?;
let service_uuid = transport.scan_for_endpoint(endpoint_id).await?;
Ok(service_uuid.to_string())
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
#[doc(hidden)]
pub async fn probe_ble_hardware_harness_native_connection(
&self,
target_node_id: &str,
) -> anyhow::Result<BleHardwareHarnessNativeProbe> {
let endpoint_id = target_node_id
.parse::<iroh::EndpointId>()
.map_err(|error| anyhow::anyhow!("invalid BLE target endpoint id: {error}"))?;
let transport = self
.ble_transport
.read()
.await
.clone()
.ok_or_else(|| anyhow::anyhow!("BLE transport has not been initialized"))?;
let report = transport.probe_endpoint_connection(endpoint_id).await;
Ok(BleHardwareHarnessNativeProbe {
endpoint_id: report.endpoint_id,
service_uuid: report.service_uuid,
device_id: report.device_id,
stages: report
.stages
.into_iter()
.map(|stage| BleHardwareHarnessNativeProbeStage {
name: stage.name,
ok: stage.ok,
elapsed_ms: stage.elapsed_ms,
detail: stage.detail,
})
.collect(),
services: report
.services
.into_iter()
.map(|service| BleHardwareHarnessNativeProbeService {
uuid: service.uuid,
characteristics: service.characteristics,
})
.collect(),
success: report.success,
error: report.error,
})
}
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(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 prune_stale_records_for_device_node(
&self,
remote_device_id: &str,
expected_node_id: &str,
) {
let records = self
.connection_manager
.get_by_device_id(remote_device_id)
.await;
let mut removed = 0usize;
let mut preserved_healthy = 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;
}
let transport_alive =
if let Ok(endpoint_id) = record_node_id.parse::<iroh::EndpointId>() {
self.is_connection_transport_alive(endpoint_id).await
} else {
false
};
if matches!(
record.state,
crate::connection_manager::ConnectionState::Connected
) && transport_alive
{
preserved_healthy = preserved_healthy.saturating_add(1);
continue;
}
self.connection_manager
.set_closing(
&record.connection_id,
Some("stale-device-node-record".to_string()),
)
.await;
self.connection_manager
.set_closed(
&record.connection_id,
Some("stale-device-node-record".to_string()),
)
.await;
let _ = self.connection_manager.remove(&record.connection_id).await;
removed = removed.saturating_add(1);
}
if removed > 0 || preserved_healthy > 0 {
let msg = format!(
"[pluto-rtc][auto-connect] pruned stale records remote_device_id={} expected_node_id={} removed={} preserved_healthy={}",
remote_device_id, expected_node_id, removed, preserved_healthy
);
#[cfg(not(target_arch = "wasm32"))]
println!("{}", msg);
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&msg));
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn retire_conflicting_device_node_records(
&self,
remote_device_id: &str,
expected_node_id: &str,
) {
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);
} else {
self.connection_manager
.set_closing(
&record.connection_id,
Some("identity-changed-retired-stale-record".to_string()),
)
.await;
self.connection_manager
.set_closed(
&record.connection_id,
Some("identity-changed-retired-stale-record".to_string()),
)
.await;
let _ = self.connection_manager.remove(&record.connection_id).await;
}
retired = retired.saturating_add(1);
}
if retired > 0 {
println!(
"[pluto-rtc][auto-connect] retired conflicting records after identity change remote_device_id={} expected_node_id={} retired={} disconnected_live={}",
remote_device_id, expected_node_id, retired, disconnected
);
}
}
#[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 ticket = match self.endpoint_ticket_with_token("user-device", 0).await {
Ok(ticket) => ticket,
Err(error) => {
eprintln!(
"[pluto-rtc][auto-connect] presence republish skipped user_id={} local_device_id={} reason={} error={}",
user_id, local_device_id, reason, error
);
return;
}
};
let (iroh_fingerprint, scope, token_fingerprint) =
summarize_compound_ticket_for_logs(ticket.as_str());
println!(
"[pluto-rtc][auto-connect][presence-republish] user_id={} local_device_id={} reason={} scope={} token_fp={} iroh_fp={}",
user_id,
local_device_id,
reason,
scope.unwrap_or_else(|| "unrestricted".to_string()),
token_fingerprint.unwrap_or_else(|| "none".to_string()),
iroh_fingerprint
);
let device_name = match self.get_native_device_identity().await {
Ok(identity) => identity.device_name,
Err(_) => crate::native_device::default_device_name().to_string(),
};
let metadata = serde_json::json!({
"deviceId": local_device_id,
"source": "auto-connect-recovery",
"reason": reason,
})
.to_string();
match self
.update_presence(user_id, &device_name, &ticket, Some(metadata.as_str()))
.await
{
Ok(()) => {
}
Err(error) => {
eprintln!(
"[pluto-rtc][auto-connect] presence republish failed user_id={} local_device_id={} reason={} error={}",
user_id, local_device_id, reason, error
);
}
}
}
#[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 known_device_id = self
.known_remote_device_id_for_incoming_transport(&connection_id, &remote_node_id)
.await;
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.connection_manager
.upsert_pending(
connection_id.clone(),
Some(remote_node_id.clone()),
known_device_id.clone(),
Some(remote_node_id.clone()),
)
.await;
self.connection_manager
.set_connected_with_transport(
&connection_id,
Some(remote_node_id.clone()),
Some(stable_id),
Some("incoming".to_string()),
)
.await;
if let Some(device_id) = known_device_id.as_deref() {
self.mark_trusted_user_device_connection_admitted(&connection_id, device_id)
.await;
}
let _ = self
.confirm_managed_connection_readiness(&connection_id)
.await;
if !already_managed_connected {
#[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 timeout_connection = connection.clone();
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 matches!(
client.session_admission(&connection_id_for_timeout),
crate::session_token::SessionAdmission::Pending
) {
#[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 {
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;
}
}
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"))]
if !already_managed_connected {
self.spawn_iroh_path_watcher(&connection_id, &remote_node_id, &connection, true);
}
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 result = node.accept_external_connection(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).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 {
self.suppress_auto_connect_for_connection_peer(
&connection_id,
"remote manual disconnect notice",
)
.await;
self.connection_manager
.clear_transport_health(&connection_id)
.await;
self.retire_managed_connection(
&connection_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
} else if close_reason.is_none() {
self.connection_manager
.set_closed(&connection_id, None)
.await;
} 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
);
self.connection_manager
.set_failed(&connection_id, Some(error.to_string()))
.await;
}
}
result
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn reconcile_incoming_transport_after_local_close(
&self,
connection_id: &str,
remote_endpoint_id: iroh::EndpointId,
remote_node_id: &str,
closed_stable_id: u64,
close_reason_debug: &str,
) -> IncomingTransportCloseResolution {
if is_manual_disconnect_close_reason(Some(close_reason_debug)) {
self.suppress_auto_connect_for_connection_peer(
connection_id,
"remote manual disconnect",
)
.await;
self.connection_manager
.clear_transport_health(connection_id)
.await;
self.retire_managed_connection(
connection_id,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
)
.await;
return IncomingTransportCloseResolution::RetiredClosedTransport;
}
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 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,
);
self.connection_manager
.set_connected_with_transport(
connection_id,
Some(remote_node_id.to_string()),
Some(kept_stable_id),
Some("incoming-kept-existing".to_string()),
)
.await;
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(
connection_id,
None,
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,
);
self.retire_managed_connection(
connection_id,
Some(
crate::lifecycle_reason::REASON_INCOMING_TRANSPORT_CLOSED_WITHOUT_LIVE_REPLACEMENT
.to_string(),
),
)
.await;
IncomingTransportCloseResolution::RetiredClosedTransport
}
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_auto_connect_excluded_peer(&device_id, node_alias, true);
#[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,
};
#[cfg(feature = "transport-lan")]
{
return crate::local_discovery::classify_iroh_path_kind(&connection);
}
#[cfg(not(feature = "transport-lan"))]
{
#[cfg(openrtc_unpublished_ble)]
{
if crate::local_discovery::selected_path_is_ble(&connection) {
return IrohPathKind::Ble;
}
}
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"),
any(feature = "transport-lan", openrtc_unpublished_ble)
))]
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"),
any(feature = "transport-lan", openrtc_unpublished_ble)
))]
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"), openrtc_unpublished_ble))]
async fn ble_transport_config(&self) -> Option<crate::client::BleConfig> {
if self.relay_only_mode_enabled().await {
return None;
}
self.transport_config
.read()
.await
.ble
.as_ref()
.filter(|ble| ble.enabled)
.cloned()
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
async fn ble_transport_enabled(&self) -> bool {
self.ble_transport_config().await.is_some()
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
async fn build_ble_transport(
endpoint_id: iroh::EndpointId,
config: crate::client::BleConfig,
) -> anyhow::Result<Arc<iroh_ble_transport::BleTransport>> {
let connect_timeout_ms = config
.connect_timeout_ms
.unwrap_or(20_000)
.clamp(500, 120_000);
let central_config = iroh_ble_transport::CentralConfig {
connect_timeout: Some(std::time::Duration::from_millis(connect_timeout_ms)),
..Default::default()
};
let central = Arc::new(iroh_ble_transport::Central::with_config(central_config).await?);
let peripheral = Arc::new(iroh_ble_transport::Peripheral::new().await?);
let retry_attempts = config.retry_attempts.unwrap_or(15).clamp(1, 100) as u32;
let retry_backoff_ms = config.retry_backoff_ms.unwrap_or(500).clamp(50, 30_000);
iroh_ble_transport::BleTransport::builder()
.l2cap_policy(iroh_ble_transport::L2capPolicy::PreferL2cap)
.retry_config(iroh_ble_transport::RetryConfig {
max_connect_attempts: retry_attempts,
base_backoff: std::time::Duration::from_millis(retry_backoff_ms),
max_backoff: std::time::Duration::from_secs(30),
})
.central(central)
.peripheral(peripheral)
.build(endpoint_id)
.await
.map_err(|error| anyhow::anyhow!("failed to initialize BLE transport: {error}"))
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
async fn apply_optional_ble_transport(
&self,
builder: iroh::endpoint::Builder,
endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::Builder> {
let Some(config) = self.ble_transport_config().await else {
return Ok(builder);
};
let ble = Self::build_ble_transport(endpoint_id, config).await?;
*self.ble_transport.write().await = Some(ble.clone());
Ok(builder
.hooks(ble.dedup_hook())
.add_custom_transport(ble.as_custom_transport())
.address_lookup(ble.address_lookup()))
}
#[cfg(all(not(target_arch = "wasm32"), not(openrtc_unpublished_ble)))]
async fn apply_optional_ble_transport(
&self,
builder: iroh::endpoint::Builder,
_endpoint_id: iroh::EndpointId,
) -> anyhow::Result<iroh::endpoint::Builder> {
Ok(builder)
}
#[cfg(all(not(target_arch = "wasm32"), openrtc_unpublished_ble))]
async fn start_ble_discovery_task(&self, node_id: &str) {
if !self.ble_transport_enabled().await {
return;
}
println!(
"[pluto-rtc][ble] iroh BLE custom transport active node_id={}",
node_id
);
}
#[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"),
any(feature = "transport-lan", openrtc_unpublished_ble)
))]
async fn start_local_discovery_tasks(&self, endpoint: &iroh::Endpoint, node_id: &str) {
#[cfg(feature = "transport-lan")]
self.start_mdns_discovery_tasks(endpoint, node_id).await;
#[cfg(not(feature = "transport-lan"))]
let _ = endpoint;
#[cfg(openrtc_unpublished_ble)]
self.start_ble_discovery_task(node_id).await;
}
#[cfg(all(
not(target_arch = "wasm32"),
not(any(feature = "transport-lan", openrtc_unpublished_ble))
))]
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);
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 == 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()),
stable_id,
Some(source.to_string()),
)
.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 } => {
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,
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.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(
&connection_id,
None,
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_manual_disconnect_close_reason(error.as_deref()) {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] remote manual disconnect — closing without replacement wait endpoint_id={} connection_id={} error={:?}",
remote_node_id, connection_id, error
)));
self.suppress_auto_connect_for_connection_peer(
&connection_id,
"remote manual disconnect",
)
.await;
self.connection_manager
.clear_transport_health(&connection_id)
.await;
self.connection_manager
.set_closed(
&connection_id,
Some(error.unwrap_or_else(|| {
crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()
})),
)
.await;
self.emit_current_wasm_connection_state(&connection_id)
.await;
return;
}
if let Some(rebound_stable_id) = self
.await_wasm_replacement_transport(
endpoint_id,
"wasm-accept-close-grace",
crate::runtime_policy::WASM_ACCEPT_CLOSE_GRACE_MS,
)
.await
{
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-accept] accepted transport close settled onto replacement endpoint_id={} connection_id={} rebound_stable_id={:?}",
remote_node_id,
connection_id,
rebound_stable_id
)));
return;
}
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(
&connection_id,
None,
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);
self.connection_manager
.mark_transport_replaced(
&connection_id,
new_stable_id,
Some("wasm-accept-close-rebind".to_string()),
None,
)
.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 {
if self.is_connected(endpoint_id).await {
let _ = self
.promote_wasm_incoming_transport(endpoint_id, source)
.await;
let rebound_stable_id = self
.get_connection(endpoint_id)
.await
.map(|connection| connection.stable_id() as u64);
if rebound_stable_id.is_some() {
return 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 graceful_disconnect_close =
is_graceful_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);
if live_transport {
web_sys::console::log_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] observed closed outgoing transport but live replacement exists endpoint_id={} closed_transport_stable_id={:?} current_live_stable_id={:?} error={:?}",
endpoint_id_str,
closed_transport_stable_id,
current_live_stable_id,
error
)));
let _ = self
.promote_wasm_incoming_transport(
endpoint_id,
"wasm-connect-bridge-rebind",
)
.await;
continue;
}
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
)));
self.connection_manager
.mark_transport_replaced(
&dup_record.connection_id,
None,
Some("replacement-held-for-accept".to_string()),
Some(
crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS
.to_string(),
),
)
.await;
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
)));
match self
.disconnect_with_reason(
endpoint_id,
crate::lifecycle_reason::REASON_ENDPOINT_HARD_RESET_REDIAL,
)
.await
{
Ok(()) => {
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
)));
}
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 = self
.connection_manager
.current_transport_matches(
&record.connection_id,
closed_transport_stable_id,
)
.await;
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(
&record.connection_id,
Some(live_replacement_stable_id),
Some("wasm-connect-bridge-closed-rebind".to_string()),
None,
)
.await;
let _ = self
.confirm_managed_connection_readiness(&record.connection_id)
.await;
} else if graceful_disconnect_close {
let reason = error.clone().or_else(|| {
Some("wasm-connect-bridge-graceful-close".to_string())
});
if is_manual_disconnect_close_reason(error.as_deref()) {
self.suppress_auto_connect_for_connection_peer(
&record.connection_id,
"remote manual disconnect",
)
.await;
}
web_sys::console::warn_1(&wasm_bindgen::JsValue::from_str(&format!(
"[pluto-rtc][wasm-connect] graceful 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
)));
self.connection_manager
.clear_transport_health(&record.connection_id)
.await;
self.connection_manager
.set_closed(&record.connection_id, reason)
.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(
&record.connection_id,
None,
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
)));
self.retire_managed_connection(
&record.connection_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() {
node.is_connected(endpoint_id).await
} else {
false
}
}
#[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 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 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,
path_kind,
force_restart,
)
.await;
}
let mut path_events = connection.path_events();
while path_events.next().await.is_some() {
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, path_kind, false)
.await;
}
});
}
}
#[cfg(not(target_arch = "wasm32"))]
async fn handle_path_change(
client: &super::Client,
connection_id: &str,
remote_node_id: &str,
path_kind: IrohPathKind,
force_restart: bool,
) {
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 webrtc_connected = client
.get_connected_webrtc_session_for_peer(connection_id)
.await
.is_some()
|| client
.get_connected_webrtc_session_for_peer(remote_node_id)
.await
.is_some();
if webrtc_connected {
client
.report_native_webrtc_candidate_transport(connection_id, Some(remote_node_id))
.await;
} else {
let _ = client
.report_transport_status(connection_id, "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 _ = client
.report_transport_status(connection_id, path_kind.transport_label(), None)
.await;
}
}