use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use futures::StreamExt;
use openrtc::application_crypto_streams::{PeerRecvStream, PeerSendStream};
use openrtc::native_node::IncomingStreamType;
use serde::{Deserialize, Serialize};
use tauri::{Emitter, Manager, Runtime};
use tokio::sync::Mutex;
const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
const RTC_DEVICE_EVENTS: &str = "rtc-device-events";
const RTC_SESSION_EVENTS: &str = "rtc-session-events";
const CONNECTION_STATE_CHANGED_EVENT: &str = "connection-state-changed";
const PEER_BI_STREAM_INCOMING_EVENT: &str = "openrtc://peer-bi-stream/incoming";
const PEER_BI_STREAM_CHUNK_EVENT: &str = "openrtc://peer-bi-stream/chunk";
const PEER_BI_STREAM_CLOSED_EVENT: &str = "openrtc://peer-bi-stream/closed";
const DISCOVERY_COMMAND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
pub type NativeTransportInstallFuture = std::pin::Pin<
Box<
dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
+ Send,
>,
>;
#[derive(Debug, Clone)]
pub struct NativeTransportInstallContext {
pub data_dir: PathBuf,
}
pub trait NativeTransportInstaller: Send + Sync {
fn id(&self) -> &'static str;
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
fn install(
&self,
client: Arc<openrtc::client::Client>,
config: openrtc::client::TransportConfig,
context: NativeTransportInstallContext,
) -> NativeTransportInstallFuture;
}
struct InstalledNativeTransport {
client_ptr: usize,
_runtime: Box<dyn std::any::Any + Send + Sync>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct OpenRtcTauriConfig {
#[serde(default = "default_project_id")]
pub project_id: String,
#[serde(default)]
pub api_key: Option<String>,
#[serde(default)]
pub app_tag: Option<String>,
#[serde(default)]
pub space_key: Option<String>,
#[serde(default)]
pub data_dir: Option<PathBuf>,
#[serde(default)]
pub transport_config: Option<openrtc::client::TransportConfig>,
}
impl Default for OpenRtcTauriConfig {
fn default() -> Self {
Self {
project_id: default_project_id(),
api_key: None,
app_tag: None,
space_key: None,
data_dir: None,
transport_config: None,
}
}
}
impl OpenRtcTauriConfig {
pub fn from_env() -> Self {
let api_key = first_env(&[
"VITE_OPENRTC_KEY",
"VITE_OPENRTC_API_KEY",
"VITE_PLUTO_OPENRTC_API_KEY",
"OPENRTC_API_KEY",
]);
let space_key = first_env(&[
"VITE_OPENRTC_SPACE_KEY",
"VITE_OPENRTC_SPACE",
"OPENRTC_SPACE_KEY",
]);
let app_tag = first_env(&[
"VITE_PLUTO_OPENRTC_APP_TAG",
"VITE_OPENRTC_APP_TAG",
"OPENRTC_APP_TAG",
]);
let project_id = first_env(&["VITE_OPENRTC_PROJECT_ID", "OPENRTC_PROJECT_ID"])
.unwrap_or_else(default_project_id);
Self {
project_id,
api_key,
app_tag,
space_key,
data_dir: None,
transport_config: None,
}
}
pub fn app_tag(&self) -> String {
if let Some(app_tag) = self
.app_tag
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
return app_tag.to_string();
}
if let (Some(api_key), Some(space_key)) = (
self.api_key
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty()),
self.space_key
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty()),
) {
return openrtc::space_app_tag_from_keys(api_key, space_key);
}
self.api_key
.as_deref()
.map(openrtc::app_tag_from_api_key)
.unwrap_or_else(|| "app_anonymous".to_string())
}
}
fn default_project_id() -> String {
openrtc::LIVE_PROJECT_ID.to_string()
}
fn first_env(names: &[&str]) -> Option<String> {
names.iter().find_map(|name| {
std::env::var(name)
.ok()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
})
}
#[derive(Default)]
struct TokenRelayState {
auth_token: RwLock<Option<String>>,
refresh_token: RwLock<Option<String>>,
}
impl TokenRelayState {
fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
let relay = self.clone();
Box::new(move || relay.auth_token.read().ok().and_then(|guard| guard.clone()))
}
fn set(&self, auth_token: Option<String>, refresh_token: Option<String>) {
if let Ok(mut guard) = self.auth_token.write() {
*guard = normalize_token(auth_token);
}
if let Ok(mut guard) = self.refresh_token.write() {
*guard = normalize_token(refresh_token);
}
}
}
fn normalize_token(value: Option<String>) -> Option<String> {
value
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
}
pub struct OpenRtcTauriState {
client: RwLock<Arc<openrtc::client::Client>>,
config: RwLock<OpenRtcTauriConfig>,
token_relay: Arc<TokenRelayState>,
data_dir: Option<PathBuf>,
connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
presence_loop_active: Mutex<bool>,
subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
managed_session_start_guard: Mutex<()>,
managed_session_next_owner_epoch: AtomicU64,
managed_session: Mutex<Option<ManagedSessionRecord>>,
native_transport_installers: Vec<Arc<dyn NativeTransportInstaller>>,
installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
}
impl OpenRtcTauriState {
pub fn new(config: OpenRtcTauriConfig) -> Self {
openrtc::ensure_default_rustls_provider();
let token_relay = Arc::new(TokenRelayState::default());
let client = build_client(&config, &token_relay);
let data_dir = config.data_dir.clone();
Self {
client: RwLock::new(client),
config: RwLock::new(config),
token_relay,
data_dir,
connection_state_forwarder: std::sync::Mutex::new(None),
presence_loop_active: Mutex::new(false),
subscriptions: Mutex::new(HashMap::new()),
peer_bi_streams: Mutex::new(HashMap::new()),
managed_session_start_guard: Mutex::new(()),
managed_session_next_owner_epoch: AtomicU64::new(0),
managed_session: Mutex::new(None),
native_transport_installers: Vec::new(),
installed_native_transports: Mutex::new(HashMap::new()),
}
}
pub fn with_native_transport_installer(
mut self,
installer: Arc<dyn NativeTransportInstaller>,
) -> Self {
self.native_transport_installers.push(installer);
self
}
pub fn client(&self) -> Arc<openrtc::client::Client> {
self.client
.read()
.map(|guard| guard.clone())
.unwrap_or_else(|_| build_client(&OpenRtcTauriConfig::default(), &self.token_relay))
}
pub fn set_auth_token(&self, auth_token: Option<String>, refresh_token: Option<String>) {
self.token_relay.set(auth_token, refresh_token);
}
fn configure_identity(
&self,
api_key: Option<String>,
space_key: Option<String>,
app_tag: Option<String>,
) -> (Arc<openrtc::client::Client>, bool) {
let mut next_config = self
.config
.read()
.map(|guard| guard.clone())
.unwrap_or_default();
apply_identity_override(&mut next_config.api_key, api_key);
apply_identity_override(&mut next_config.space_key, space_key);
apply_identity_override(&mut next_config.app_tag, app_tag);
let next_app_tag = next_config.app_tag();
let current = self.client();
if current.app_tag() == next_app_tag {
if let Ok(mut guard) = self.config.write() {
*guard = next_config;
}
return (current, false);
}
let next_client = build_client(&next_config, &self.token_relay);
if let Ok(mut guard) = self.client.write() {
*guard = next_client.clone();
}
if let Ok(mut guard) = self.config.write() {
*guard = next_config;
}
current.stop_auth_scoped_activity();
(next_client, true)
}
fn allocate_managed_session_owner_epoch(&self) -> u64 {
self.managed_session_next_owner_epoch
.fetch_add(1, Ordering::SeqCst)
+ 1
}
fn replace_connection_state_forwarder<R: Runtime>(
&self,
app: tauri::AppHandle<R>,
client: Arc<openrtc::client::Client>,
) {
let next = tauri::async_runtime::spawn(async move {
forward_connection_state_events(app, client).await;
});
if let Ok(mut guard) = self.connection_state_forwarder.lock() {
if let Some(previous) = guard.replace(next) {
previous.abort();
}
}
}
}
fn build_client(
config: &OpenRtcTauriConfig,
token_relay: &Arc<TokenRelayState>,
) -> Arc<openrtc::client::Client> {
let mut builder = openrtc::client::Client::builder_with_app_tag(
config.project_id.clone(),
config.app_tag(),
token_relay.token_provider(),
);
if let Some(transport_config) = config.transport_config.clone() {
builder = builder.transport_config(transport_config);
}
Arc::new(builder.build())
}
fn apply_identity_override(target: &mut Option<String>, value: Option<String>) {
if let Some(value) = value
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
*target = Some(value.to_string());
}
}
fn requested_local_device_id(value: Option<&str>) -> Option<String> {
value
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn managed_session_device_id(
_requested: Option<&str>,
native_identity: &openrtc::native_device::NativeDeviceIdentity,
) -> String {
native_identity.device_id.clone()
}
fn metadata_with_authoritative_device_id(
metadata: Option<String>,
device_id: &str,
) -> Option<String> {
let device_id = device_id.trim();
if device_id.is_empty() {
return metadata;
}
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.insert(
"deviceId".to_string(),
serde_json::Value::String(device_id.to_string()),
);
Some(serde_json::Value::Object(map).to_string())
}
_ => Some(
serde_json::json!({
"deviceId": device_id,
"metadata": raw_metadata,
})
.to_string(),
),
}
}
struct PeerBiStreamHandle {
send: Arc<Mutex<Option<PeerSendStream>>>,
recv: Option<PeerRecvStream>,
read_task: Option<tokio::task::JoinHandle<()>>,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
pub struct OpenPeerBiStreamResult {
stream_id: String,
connection_id: Option<String>,
remote_node_id: String,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct IncomingPeerBiStreamEvent {
request_id: String,
stream_id: String,
remote_node_id: String,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct PeerBiStreamChunkEvent {
stream_id: String,
bytes: Vec<u8>,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct PeerBiStreamClosedEvent {
stream_id: String,
error: Option<String>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SearchRtcDevicesResult {
devices: Vec<openrtc::signaling::Device>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SearchRtcDevicesWithStatusResult {
devices: Vec<openrtc::client::DeviceStatusSnapshot>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct StartRtcManagedSessionResult {
local_node_id: String,
ticket_scope: Option<String>,
ticket: Option<String>,
presence_started: bool,
auto_connect_started: bool,
local_device: openrtc::native_device::NativeDeviceIdentity,
}
#[derive(Clone)]
struct ManagedSessionRecord {
key: String,
owner_client: Arc<openrtc::client::Client>,
owner_epoch: u64,
result: StartRtcManagedSessionResult,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ManagedSessionDisposition {
Start,
Reuse,
Refresh,
Replace,
}
fn managed_session_disposition(
active_key: Option<&str>,
active_client_matches: bool,
active_presence_started: bool,
active_auto_connect_started: bool,
requested_key: &str,
requested_presence: bool,
requested_auto_connect: bool,
) -> ManagedSessionDisposition {
let Some(active_key) = active_key else {
return ManagedSessionDisposition::Start;
};
if !active_client_matches || active_key != requested_key {
return ManagedSessionDisposition::Replace;
}
if (!requested_presence || active_presence_started)
&& (!requested_auto_connect || active_auto_connect_started)
{
ManagedSessionDisposition::Reuse
} else {
ManagedSessionDisposition::Refresh
}
}
pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
init_with_state(OpenRtcTauriState::new(config))
}
pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
tauri::plugin::Builder::new(PLUGIN_NAME)
.setup(move |app, _api| {
let client = state.client();
let transport_config = state
.config
.read()
.ok()
.and_then(|config| config.transport_config.clone());
let install_context = NativeTransportInstallContext {
data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
};
tauri::async_runtime::block_on(ensure_requested_native_transports(
&state,
&client,
transport_config.as_ref(),
&install_context,
))
.map_err(std::io::Error::other)?;
state.replace_connection_state_forwarder(app.clone(), client);
app.manage(state);
Ok(())
})
.invoke_handler(tauri::generate_handler![
desktop_set_pluto_auth_token,
openrtc_set_auth_token,
rtc_native_status,
get_rtc_local_device_info,
update_rtc_local_device_name,
get_iroh_node_id,
start_iroh_node,
get_iroh_endpoint_ticket,
register_session_token,
get_endpoint_ticket_with_token,
validate_session_token,
revoke_session_tokens_by_scope,
search_rtc_devices,
search_rtc_devices_with_status,
update_rtc_presence,
update_rtc_device,
delete_rtc_device,
set_rtc_offline,
start_rtc_presence_loop,
stop_rtc_presence_loop,
start_rtc_managed_session,
start_rtc_auto_connect,
stop_rtc_auto_connect,
force_rtc_reconnect_snapshot,
notify_rtc_network_change,
connect_to_device,
disconnect_device,
set_auto_connect_excluded,
resolve_rtc_peer_connection_records,
resolve_rtc_peer_identity,
get_rtc_peer_session,
list_rtc_peer_sessions,
list_rtc_managed_connections,
wait_for_rtc_settled_peer,
list_rtc_connection_states,
get_rtc_connection_state,
start_rtc_device_subscription,
start_rtc_session_subscription,
stop_rtc_subscription,
start_incoming_peer_bi_streams,
open_peer_bi_stream,
open_peer_bi_transport_only_stream,
open_peer_native_bi_stream,
write_peer_bi_stream,
start_peer_bi_stream_read,
close_peer_bi_stream,
get_app_limits,
])
.build()
}
pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
init(OpenRtcTauriConfig::from_env())
}
async fn forward_connection_state_events<R: Runtime>(
app: tauri::AppHandle<R>,
client: Arc<openrtc::client::Client>,
) {
let mut rx = client.subscribe_native_connection_state_updates();
while let Ok(snapshot) = rx.recv().await {
let _ = app.emit(CONNECTION_STATE_CHANGED_EVENT, snapshot);
}
}
fn app_data_dir<R: Runtime>(
app: &tauri::AppHandle<R>,
state: &OpenRtcTauriState,
) -> Result<PathBuf, String> {
state
.data_dir
.clone()
.or_else(|| app.path().app_data_dir().ok())
.ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
}
async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
if let Some(node_id) = client.current_node_id().await {
return Ok(node_id);
}
client
.init_iroh(None, Vec::new())
.await
.map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
}
async fn ensure_requested_native_transports(
state: &OpenRtcTauriState,
client: &Arc<openrtc::client::Client>,
transports: Option<&openrtc::client::TransportConfig>,
context: &NativeTransportInstallContext,
) -> Result<(), String> {
let Some(transports) = transports else {
return Ok(());
};
let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
if ble_requested
&& !state
.native_transport_installers
.iter()
.any(|installer| installer.is_requested(transports))
{
eprintln!(
"[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
);
return Ok(());
}
let client_ptr = Arc::as_ptr(client) as usize;
for installer in &state.native_transport_installers {
if !installer.is_requested(transports) {
continue;
}
let already_installed = state
.installed_native_transports
.lock()
.await
.get(installer.id())
.is_some_and(|installed| installed.client_ptr == client_ptr);
if already_installed {
continue;
}
if client.current_node_id().await.is_some() {
eprintln!(
"[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
installer.id()
);
continue;
}
let runtime = match installer
.install(client.clone(), transports.clone(), context.clone())
.await
{
Ok(runtime) => runtime,
Err(error) => {
eprintln!(
"[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
installer.id(),
error
);
continue;
}
};
state.installed_native_transports.lock().await.insert(
installer.id(),
InstalledNativeTransport {
client_ptr,
_runtime: runtime,
},
);
}
Ok(())
}
async fn register_peer_bi_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: &OpenRtcTauriState,
connection_id: Option<String>,
remote_node_id: String,
send: PeerSendStream,
recv: PeerRecvStream,
start_reading: bool,
) -> OpenPeerBiStreamResult {
let stream_id = uuid::Uuid::new_v4().to_string();
let send = Arc::new(Mutex::new(Some(send)));
let (recv, read_task) = if start_reading {
(
None,
Some(spawn_peer_bi_stream_reader(
app.clone(),
stream_id.clone(),
recv,
)),
)
} else {
(Some(recv), None)
};
state.peer_bi_streams.lock().await.insert(
stream_id.clone(),
PeerBiStreamHandle {
send,
recv,
read_task,
},
);
OpenPeerBiStreamResult {
stream_id,
connection_id,
remote_node_id,
}
}
fn spawn_peer_bi_stream_reader<R: Runtime>(
app: tauri::AppHandle<R>,
stream_id: String,
mut recv: PeerRecvStream,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut chunk = vec![0_u8; 64 * 1024];
let mut close_error: Option<String> = None;
loop {
match recv.read(&mut chunk).await {
Ok(0) => break,
Ok(n) => {
let _ = app.emit(
PEER_BI_STREAM_CHUNK_EVENT,
PeerBiStreamChunkEvent {
stream_id: stream_id.clone(),
bytes: chunk[..n].to_vec(),
},
);
}
Err(error) => {
close_error = Some(error.to_string());
break;
}
}
}
let _ = app.emit(
PEER_BI_STREAM_CLOSED_EVENT,
PeerBiStreamClosedEvent {
stream_id,
error: close_error,
},
);
})
}
async fn open_peer_bi_with<R: Runtime, F, Fut>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
open: F,
) -> Result<OpenPeerBiStreamResult, String>
where
F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
Fut: std::future::Future<
Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
>,
{
let peer_id = peer_id.trim().to_string();
if peer_id.is_empty() {
return Err("peerId is required".to_string());
}
let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
.await
.map_err(|error| format!("open peer bi stream failed: {error}"))?;
Ok(register_peer_bi_stream(app, &state, connection_id, remote_node_id, send, recv, true).await)
}
#[tauri::command]
async fn desktop_set_pluto_auth_token(
state: tauri::State<'_, OpenRtcTauriState>,
auth_token: Option<String>,
refresh_token: Option<String>,
) -> Result<(), String> {
state.set_auth_token(auth_token, refresh_token);
if *state.presence_loop_active.lock().await {
state.client().request_presence_update();
}
Ok(())
}
#[tauri::command]
async fn openrtc_set_auth_token(
state: tauri::State<'_, OpenRtcTauriState>,
auth_token: Option<String>,
refresh_token: Option<String>,
) -> Result<(), String> {
desktop_set_pluto_auth_token(state, auth_token, refresh_token).await
}
#[tauri::command]
async fn rtc_native_status(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<openrtc::client::RuntimeStatus, String> {
Ok(state.client().runtime_status().await)
}
#[tauri::command]
async fn get_rtc_local_device_info<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
let data_dir = app_data_dir(&app, &state)?;
state
.client()
.init_native_device_identity(data_dir, None)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn update_rtc_local_device_name<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
device_name: String,
) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
let data_dir = app_data_dir(&app, &state)?;
let _ = state
.client()
.init_native_device_identity(data_dir, None)
.await
.map_err(|error| error.to_string())?;
state
.client()
.update_native_device_name(&device_name)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn get_iroh_node_id(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Option<String>, String> {
Ok(state.client().current_node_id().await)
}
#[tauri::command]
async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
ensure_iroh_node(state.client()).await
}
#[tauri::command]
async fn get_iroh_endpoint_ticket(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<String, String> {
ensure_iroh_node(state.client()).await?;
state
.client()
.endpoint_ticket()
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn register_session_token(
state: tauri::State<'_, OpenRtcTauriState>,
token: String,
scope: String,
max_connections: u32,
expires_at_ms: Option<u64>,
) -> Result<(), String> {
if let Some(expires_at_ms) = expires_at_ms {
state.client().register_session_token_with_expiry_ms(
token,
scope,
max_connections,
expires_at_ms,
);
} else {
state
.client()
.register_session_token(token, scope, max_connections);
}
Ok(())
}
#[tauri::command]
async fn get_endpoint_ticket_with_token(
state: tauri::State<'_, OpenRtcTauriState>,
scope: String,
max_connections: u32,
) -> Result<String, String> {
ensure_iroh_node(state.client()).await?;
state
.client()
.endpoint_ticket_with_token(&scope, max_connections)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn validate_session_token(
state: tauri::State<'_, OpenRtcTauriState>,
token: String,
connection_id: Option<String>,
) -> Result<String, String> {
if let Some(connection_id) = connection_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
state
.client()
.validate_session_token_for_connection(&token, connection_id)
.await
.map_err(|error| error.to_string())
} else {
state.client().validate_session_token(&token)
}
}
#[tauri::command]
async fn revoke_session_tokens_by_scope(
state: tauri::State<'_, OpenRtcTauriState>,
scope: String,
) -> Result<Vec<String>, String> {
Ok(state.client().revoke_tokens_by_scope(&scope).await)
}
fn discovery_verbose() -> bool {
std::env::var("OPENRTC_SIGNALING_VERBOSE")
.map(|value| {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"1" | "true" | "yes" | "on"
)
})
.unwrap_or(false)
}
#[tauri::command]
async fn search_rtc_devices(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
) -> Result<SearchRtcDevicesResult, String> {
let client = state.client();
let app_tag = client.app_tag().to_string();
let started = std::time::Instant::now();
if discovery_verbose() {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices start user_id={} app_tag={}",
user_id, app_tag
);
}
let devices = match tokio::time::timeout(
DISCOVERY_COMMAND_TIMEOUT,
client.search_devices(&user_id),
)
.await
{
Ok(Ok(devices)) => devices,
Ok(Err(error)) => {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices failed user_id={} app_tag={} elapsed_ms={} error={}",
user_id,
app_tag,
started.elapsed().as_millis(),
error
);
return Err(error.to_string());
}
Err(_) => {
let message = format!(
"search_rtc_devices native discovery timed out after {}ms user_id={} app_tag={}",
DISCOVERY_COMMAND_TIMEOUT.as_millis(),
user_id,
app_tag
);
eprintln!(
"[openrtc-tauri][discovery] {} elapsed_ms={}",
message,
started.elapsed().as_millis()
);
return Err(message);
}
};
if discovery_verbose() {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices done user_id={} app_tag={} count={} elapsed_ms={}",
user_id,
app_tag,
devices.len(),
started.elapsed().as_millis()
);
}
Ok(SearchRtcDevicesResult { devices })
}
#[tauri::command]
async fn search_rtc_devices_with_status(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
) -> Result<SearchRtcDevicesWithStatusResult, String> {
let client = state.client();
let app_tag = client.app_tag().to_string();
let started = std::time::Instant::now();
if discovery_verbose() {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices_with_status start user_id={} app_tag={}",
user_id, app_tag
);
}
let devices = match tokio::time::timeout(
DISCOVERY_COMMAND_TIMEOUT,
client.devices_with_status(&user_id),
)
.await
{
Ok(Ok(devices)) => devices,
Ok(Err(error)) => {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices_with_status failed user_id={} app_tag={} elapsed_ms={} error={}",
user_id,
app_tag,
started.elapsed().as_millis(),
error
);
return Err(error.to_string());
}
Err(_) => {
let message = format!(
"search_rtc_devices_with_status native discovery timed out after {}ms user_id={} app_tag={}",
DISCOVERY_COMMAND_TIMEOUT.as_millis(),
user_id,
app_tag
);
eprintln!(
"[openrtc-tauri][discovery] {} elapsed_ms={}",
message,
started.elapsed().as_millis()
);
return Err(message);
}
};
if discovery_verbose() {
eprintln!(
"[openrtc-tauri][discovery] search_rtc_devices_with_status done user_id={} app_tag={} count={} elapsed_ms={}",
user_id,
app_tag,
devices.len(),
started.elapsed().as_millis()
);
}
Ok(SearchRtcDevicesWithStatusResult { devices })
}
#[tauri::command]
async fn update_rtc_presence(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
) -> Result<(), String> {
state
.client()
.update_durable_device_record_with_ttl(
&user_id,
&device_name,
&ticket,
7 * 24 * 60 * 60 * 1000,
metadata.as_deref(),
)
.await
.map_err(|error| error.to_string())?;
state
.client()
.update_live_presence_record(&user_id, &device_name, &ticket, metadata.as_deref())
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn update_rtc_device(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
device_id: String,
device_name: Option<String>,
capabilities: Option<openrtc::signaling::DeviceCapabilities>,
metadata: Option<String>,
) -> Result<(), String> {
state
.client()
.update_device(
&user_id,
&device_id,
device_name.as_deref(),
capabilities,
metadata.as_deref(),
)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn delete_rtc_device(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
device_id: String,
) -> Result<(), String> {
state
.client()
.delete_device(&user_id, &device_id)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn set_rtc_offline(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
) -> Result<(), String> {
state
.client()
.set_offline(&user_id)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn start_rtc_presence_loop(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
local_node_id: String,
device_name: String,
ticket: String,
metadata: Option<String>,
api_key: Option<String>,
space_key: Option<String>,
) -> Result<(), String> {
let _ = local_node_id;
let (client, _) = state.configure_identity(api_key.clone(), space_key, None);
if let Some(api_key) = api_key
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
let api_key = api_key.to_string();
let client = client.clone();
tokio::spawn(async move {
let _ = client.fetch_and_cache_app_limits(&api_key).await;
});
}
client.start_signaling_loop(
user_id.clone(),
device_name.clone(),
ticket.clone(),
metadata.clone(),
);
*state.presence_loop_active.lock().await = true;
Ok(())
}
#[tauri::command]
async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
state.client().stop_presence_loop();
*state.presence_loop_active.lock().await = false;
if let Some(active) = state.managed_session.lock().await.as_mut() {
active.result.presence_started = false;
}
stop_subscription_by_id(&state, "internal-presence-loop").await;
Ok(())
}
#[tauri::command]
async fn start_rtc_managed_session<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
device_name: Option<String>,
local_device_id: Option<String>,
metadata: Option<String>,
transports: Option<openrtc::client::TransportConfig>,
auto_connect: Option<bool>,
presence: Option<bool>,
api_key: Option<String>,
space_key: Option<String>,
) -> Result<StartRtcManagedSessionResult, String> {
let command_started = std::time::Instant::now();
let _start_guard = state.managed_session_start_guard.lock().await;
let (client, reconfigured) = state.configure_identity(api_key.clone(), space_key.clone(), None);
let app_tag = client.app_tag().to_string();
eprintln!(
"[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
user_id,
app_tag,
device_name
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
local_device_id
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
metadata
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
auto_connect.unwrap_or(true),
presence.unwrap_or(true)
);
if reconfigured {
state.replace_connection_state_forwarder(app.clone(), client.clone());
eprintln!(
"[openrtc-tauri] reconfigured native OpenRTC namespace app_tag={}",
client.app_tag()
);
}
let data_dir = app_data_dir(&app, &state)?;
let install_context = NativeTransportInstallContext {
data_dir: data_dir.clone(),
};
ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
.await?;
if let Some(transport_config) = transports {
client
.update_transport_config(transport_config)
.await
.map_err(|error| error.to_string())?;
}
let local_node_id = ensure_iroh_node(client.clone()).await?;
let local_device = client
.init_native_device_identity(data_dir, device_name.as_deref())
.await
.map_err(|error| error.to_string())?;
let effective_local_device_id =
managed_session_device_id(local_device_id.as_deref(), &local_device);
if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
if requested != local_device.device_id {
eprintln!(
"[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
requested, local_device.device_id
);
}
}
let should_start_presence = presence.unwrap_or(true);
let should_start_auto_connect = auto_connect.unwrap_or(true);
let managed_session_key = format!("{}:{}:{}", app_tag, user_id, effective_local_device_id);
let (disposition, reused) = {
let active = state.managed_session.lock().await;
let disposition = managed_session_disposition(
active.as_ref().map(|record| record.key.as_str()),
active
.as_ref()
.is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
active
.as_ref()
.is_some_and(|record| record.result.presence_started),
active
.as_ref()
.is_some_and(|record| record.result.auto_connect_started),
&managed_session_key,
should_start_presence,
should_start_auto_connect,
);
let reused = (disposition == ManagedSessionDisposition::Reuse).then(|| {
active
.as_ref()
.expect("managed session reuse requires an active record")
.result
.clone()
});
(disposition, reused)
};
if let Some(reused) = reused {
eprintln!(
"[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
managed_session_key,
command_started.elapsed().as_millis()
);
return Ok(reused);
}
if disposition == ManagedSessionDisposition::Replace {
let previous = state
.managed_session
.lock()
.await
.take()
.expect("managed session replacement requires an active record");
previous.owner_client.stop_presence_loop();
previous.owner_client.stop_auto_connect();
*state.presence_loop_active.lock().await = false;
}
let ticket_scope = should_start_presence.then(|| "user-device".to_string());
let managed_metadata =
metadata_with_authoritative_device_id(metadata.clone(), &effective_local_device_id);
if should_start_presence {
client
.clone()
.start_managed_user_device_presence_loop(
user_id.clone(),
local_device.device_name.clone(),
managed_metadata.clone(),
)
.await
.map_err(|error| {
format!(
"managed native presence failed initial publication: {}",
error
)
})?;
*state.presence_loop_active.lock().await = true;
if let Some(api_key) = api_key
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
let api_key = api_key.to_string();
let limits_client = client.clone();
tokio::spawn(async move {
let _ = limits_client.fetch_and_cache_app_limits(&api_key).await;
});
}
}
if should_start_auto_connect {
client
.clone()
.start_auto_connect(user_id.clone(), effective_local_device_id.clone());
}
let logical_local_device = local_device.clone();
let owner_epoch = match disposition {
ManagedSessionDisposition::Refresh => {
state
.managed_session
.lock()
.await
.as_ref()
.expect("managed session refresh requires an active record")
.owner_epoch
}
ManagedSessionDisposition::Start | ManagedSessionDisposition::Replace => {
state.allocate_managed_session_owner_epoch()
}
ManagedSessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
};
eprintln!(
"[openrtc-tauri][managed-session] ready user_id={} app_tag={} local_node_id={} local_device_id={} owner_epoch={} presence_started={} auto_connect_started={} ticket_scope={} elapsed_ms={}",
user_id,
app_tag,
local_node_id,
logical_local_device.device_id,
owner_epoch,
should_start_presence,
should_start_auto_connect,
ticket_scope.as_deref().unwrap_or("unrestricted"),
command_started.elapsed().as_millis()
);
let result = StartRtcManagedSessionResult {
local_node_id,
ticket_scope,
ticket: None,
presence_started: should_start_presence,
auto_connect_started: should_start_auto_connect,
local_device: logical_local_device,
};
*state.managed_session.lock().await = Some(ManagedSessionRecord {
key: managed_session_key,
owner_client: client,
owner_epoch,
result: result.clone(),
});
Ok(result)
}
#[tauri::command]
async fn start_rtc_auto_connect(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
local_device_id: String,
) -> Result<(), String> {
state
.client()
.clone()
.start_auto_connect(user_id, local_device_id);
Ok(())
}
#[tauri::command]
async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
state.client().stop_auto_connect();
if let Some(active) = state.managed_session.lock().await.as_mut() {
active.result.auto_connect_started = false;
}
Ok(())
}
#[tauri::command]
async fn force_rtc_reconnect_snapshot(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<(), String> {
state.client().clone().force_reconnect_snapshot();
Ok(())
}
#[tauri::command]
async fn connect_to_device(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: Option<String>,
endpoint_ticket: String,
timeout_ms: Option<u64>,
) -> Result<openrtc::client::ManagedConnectResult, String> {
let client = state.client();
let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
match timeout_ms {
Some(timeout_ms) => {
tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
.await
.map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
.map_err(|error| error.to_string())
}
None => connect.await.map_err(|error| error.to_string()),
}
}
#[tauri::command]
async fn disconnect_device(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
node_id_hint: Option<String>,
) -> Result<(), String> {
state
.client()
.disconnect_device(&device_id, node_id_hint.as_deref())
.await;
Ok(())
}
#[tauri::command]
async fn set_auto_connect_excluded(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
excluded: bool,
) -> Result<(), String> {
if excluded {
state.client().exclude_peer_and_publish(&device_id).await;
} else {
state.client().unexclude_peer_and_publish(&device_id).await;
}
Ok(())
}
#[tauri::command]
async fn resolve_rtc_peer_connection_records(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
Ok(state.client().resolve_peer_connection_records(&id).await)
}
#[tauri::command]
async fn resolve_rtc_peer_identity(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
Ok(state.client().peer_snapshot(&id).await)
}
#[tauri::command]
async fn get_rtc_peer_session(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().peer_session(&id).await)
}
#[tauri::command]
async fn list_rtc_peer_sessions(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().peer_sessions().await)
}
#[tauri::command]
async fn list_rtc_managed_connections(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
Ok(state.client().list_managed_connections().await)
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct NotifyRtcNetworkChangeResult {
retired_stale_connections: usize,
}
#[tauri::command]
async fn notify_rtc_network_change(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<NotifyRtcNetworkChangeResult, String> {
let retired_stale_connections = state
.client()
.notify_network_change()
.await
.map_err(|error| error.to_string())?;
Ok(NotifyRtcNetworkChangeResult {
retired_stale_connections,
})
}
#[tauri::command]
async fn wait_for_rtc_settled_peer(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
timeout_ms: Option<u64>,
) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().wait_for_settled_peer(&id, timeout_ms).await)
}
#[tauri::command]
async fn list_rtc_connection_states(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Vec<openrtc::client::ConnectionStateSnapshot>, String> {
Ok(state.client().connection_states().await)
}
#[tauri::command]
async fn get_rtc_connection_state(
state: tauri::State<'_, OpenRtcTauriState>,
connection_id: String,
) -> Result<Option<openrtc::client::ConnectionStateSnapshot>, String> {
Ok(state.client().connection_state(&connection_id).await)
}
#[tauri::command]
async fn start_rtc_device_subscription<R: Runtime>(
state: tauri::State<'_, OpenRtcTauriState>,
app: tauri::AppHandle<R>,
request_id: String,
user_id: String,
) -> Result<(), String> {
stop_subscription_by_id(&state, &request_id).await;
let mut stream = state
.client()
.subscribe_devices(&user_id)
.await
.map_err(|error| error.to_string())?;
let request_for_task = request_id.clone();
let handle = tokio::spawn(async move {
while let Some(result) = stream.next().await {
let payload = match result {
Ok(events) => serde_json::json!({
"requestId": request_for_task,
"type": "deviceEvents",
"events": events,
}),
Err(error) => serde_json::json!({
"requestId": request_for_task,
"type": "error",
"message": error.to_string(),
}),
};
let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
let _ = app.emit(RTC_DEVICE_EVENTS, payload);
if is_error {
break;
}
}
});
state.subscriptions.lock().await.insert(request_id, handle);
Ok(())
}
#[tauri::command]
async fn start_rtc_session_subscription<R: Runtime>(
state: tauri::State<'_, OpenRtcTauriState>,
app: tauri::AppHandle<R>,
request_id: String,
local_device_id: String,
) -> Result<(), String> {
stop_subscription_by_id(&state, &request_id).await;
let mut stream = state
.client()
.subscribe_sessions(&local_device_id)
.await
.map_err(|error| error.to_string())?;
let request_for_task = request_id.clone();
let handle = tokio::spawn(async move {
while let Some(result) = stream.next().await {
let payload = match result {
Ok(events) => serde_json::json!({
"requestId": request_for_task,
"type": "sessionEvents",
"events": events,
}),
Err(error) => serde_json::json!({
"requestId": request_for_task,
"type": "error",
"message": error.to_string(),
}),
};
let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
let _ = app.emit(RTC_SESSION_EVENTS, payload);
if is_error {
break;
}
}
});
state.subscriptions.lock().await.insert(request_id, handle);
Ok(())
}
#[tauri::command]
async fn stop_rtc_subscription(
state: tauri::State<'_, OpenRtcTauriState>,
request_id: String,
) -> Result<(), String> {
stop_subscription_by_id(&state, &request_id).await;
Ok(())
}
async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
handle.abort();
}
}
#[tauri::command]
async fn start_incoming_peer_bi_streams<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
request_id: Option<String>,
) -> Result<String, String> {
let request_id = request_id
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
stop_subscription_by_id(&state, &request_id).await;
let incoming = state
.client()
.incoming_streams()
.await
.map_err(|error| format!("failed to subscribe to incoming peer streams: {error}"))?;
let app_for_task = app.clone();
let request_for_task = request_id.clone();
let handle = tokio::spawn(async move {
while let Ok(incoming_stream) = incoming.recv().await {
let remote_node_id = incoming_stream.endpoint_id.to_string();
let IncomingStreamType::Bi(send, recv) = incoming_stream.stream else {
continue;
};
let state_for_task = app_for_task.state::<OpenRtcTauriState>();
let result = register_peer_bi_stream(
app_for_task.clone(),
state_for_task.inner(),
None,
remote_node_id.clone(),
PeerSendStream::plain(send),
PeerRecvStream::plain(recv),
false,
)
.await;
let _ = app_for_task.emit(
PEER_BI_STREAM_INCOMING_EVENT,
IncomingPeerBiStreamEvent {
request_id: request_for_task.clone(),
stream_id: result.stream_id,
remote_node_id,
},
);
}
});
state
.subscriptions
.lock()
.await
.insert(request_id.clone(), handle);
Ok(request_id)
}
#[tauri::command]
async fn open_peer_bi_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
) -> Result<OpenPeerBiStreamResult, String> {
open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
client.open_peer_bi(&peer_id, timeout_ms).await
})
.await
}
#[tauri::command]
async fn open_peer_bi_transport_only_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
) -> Result<OpenPeerBiStreamResult, String> {
open_peer_bi_with(
app,
state,
peer_id,
timeout_ms,
|client, peer_id, timeout_ms| async move {
client
.open_peer_bi_transport_only(&peer_id, timeout_ms)
.await
.map(|(connection_id, remote_node_id, send, recv)| {
(
connection_id,
remote_node_id,
PeerSendStream::plain(send),
PeerRecvStream::plain(recv),
)
})
},
)
.await
}
#[tauri::command]
async fn open_peer_native_bi_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
label: String,
timeout_ms: Option<u64>,
) -> Result<OpenPeerBiStreamResult, String> {
let channel_envelope = openrtc::stream_metadata::encode_channel_envelope(&label, None)
.map_err(|error| error.to_string())?;
open_peer_bi_with(
app,
state,
peer_id,
timeout_ms,
move |client, peer_id, timeout_ms| async move {
let (connection_id, remote_node_id, mut send, recv) =
client.open_peer_bi(&peer_id, timeout_ms).await?;
send.write_all(&channel_envelope).await?;
Ok((connection_id, remote_node_id, send, recv))
},
)
.await
}
#[tauri::command]
async fn write_peer_bi_stream(
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
bytes: Vec<u8>,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let send = {
let streams = state.peer_bi_streams.lock().await;
streams
.get(&stream_id)
.map(|handle| handle.send.clone())
.ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
};
let result = send
.lock()
.await
.as_mut()
.ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
.write_all(&bytes)
.await
.map_err(|error| format!("write peer bi stream failed: {error}"));
result
}
#[tauri::command]
async fn start_peer_bi_stream_read<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let mut streams = state.peer_bi_streams.lock().await;
let handle = streams
.get_mut(&stream_id)
.ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
if handle.read_task.is_some() {
return Ok(());
}
let recv = handle
.recv
.take()
.ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
handle.read_task = Some(spawn_peer_bi_stream_reader(app, stream_id, recv));
Ok(())
}
#[tauri::command]
async fn close_peer_bi_stream(
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let handle = {
let mut streams = state.peer_bi_streams.lock().await;
streams.remove(&stream_id)
};
if let Some(handle) = handle {
let mut send = handle.send.lock().await;
if let Some(send) = send.take() {
let _ = send.finish();
}
drop(send);
if let Some(read_task) = handle.read_task {
read_task.abort();
}
}
Ok(())
}
#[tauri::command]
async fn get_app_limits() -> Result<serde_json::Value, String> {
Ok(serde_json::json!({
"devicesPerUser": -1,
"maxRooms": -1,
"maxMembersPerRoom": -1
}))
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::oneshot;
struct FakeNativeTransportInstaller {
installs: Arc<AtomicUsize>,
}
impl NativeTransportInstaller for FakeNativeTransportInstaller {
fn id(&self) -> &'static str {
"ble"
}
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
config.ble.as_ref().is_some_and(|ble| ble.enabled)
}
fn install(
&self,
_client: Arc<openrtc::client::Client>,
_config: openrtc::client::TransportConfig,
_context: NativeTransportInstallContext,
) -> NativeTransportInstallFuture {
self.installs.fetch_add(1, Ordering::SeqCst);
Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
}
}
#[derive(Debug)]
struct UnavailableNativeTransportInstaller;
impl NativeTransportInstaller for UnavailableNativeTransportInstaller {
fn id(&self) -> &'static str {
"ble"
}
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
config.ble.as_ref().is_some_and(|ble| ble.enabled)
}
fn install(
&self,
_client: Arc<openrtc::client::Client>,
_config: openrtc::client::TransportConfig,
_context: NativeTransportInstallContext,
) -> NativeTransportInstallFuture {
Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
}
}
fn requested_ble_config() -> openrtc::client::TransportConfig {
openrtc::client::TransportConfig {
ble: Some(openrtc::client::BleConfig {
enabled: true,
..openrtc::client::BleConfig::default()
}),
..openrtc::client::TransportConfig::default()
}
}
fn native_transport_install_context() -> NativeTransportInstallContext {
NativeTransportInstallContext {
data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
}
}
fn managed_session_test_result(user_id: &str) -> StartRtcManagedSessionResult {
StartRtcManagedSessionResult {
local_node_id: format!("node-{user_id}"),
ticket_scope: Some("user-device".to_string()),
ticket: None,
presence_started: true,
auto_connect_started: true,
local_device: openrtc::native_device::NativeDeviceIdentity {
device_id: format!("device-{user_id}"),
device_name: "Test Device".to_string(),
created_at_ms: 1,
updated_at_ms: 1,
name_source: None,
system_info: None,
},
}
}
async fn simulate_managed_session_start(
state: Arc<OpenRtcTauriState>,
api_key: String,
user_id: String,
pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
starts: Arc<AtomicUsize>,
stopped_owners: Arc<Mutex<Vec<usize>>>,
) -> ManagedSessionDisposition {
let _start_guard = state.managed_session_start_guard.lock().await;
let (client, _) = state.configure_identity(Some(api_key), None, None);
let key = format!("{}:{user_id}:device-{user_id}", client.app_tag());
let (disposition, active_epoch) = {
let active = state.managed_session.lock().await;
let disposition = managed_session_disposition(
active.as_ref().map(|record| record.key.as_str()),
active
.as_ref()
.is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
active
.as_ref()
.is_some_and(|record| record.result.presence_started),
active
.as_ref()
.is_some_and(|record| record.result.auto_connect_started),
&key,
true,
true,
);
(
disposition,
active.as_ref().map(|record| record.owner_epoch),
)
};
if disposition == ManagedSessionDisposition::Reuse {
return disposition;
}
if disposition == ManagedSessionDisposition::Replace {
let previous = state
.managed_session
.lock()
.await
.take()
.expect("replacement must have a previous owner");
stopped_owners
.lock()
.await
.push(Arc::as_ptr(&previous.owner_client) as usize);
}
starts.fetch_add(1, Ordering::SeqCst);
if let Some((entered, release)) = pause {
entered.send(()).expect("start observer must be waiting");
release.await.expect("start release must be sent");
}
let owner_epoch = match disposition {
ManagedSessionDisposition::Refresh => {
active_epoch.expect("refresh must preserve the active owner epoch")
}
ManagedSessionDisposition::Start | ManagedSessionDisposition::Replace => {
state.allocate_managed_session_owner_epoch()
}
ManagedSessionDisposition::Reuse => unreachable!("reuse returned before startup"),
};
let result = managed_session_test_result(&user_id);
*state.managed_session.lock().await = Some(ManagedSessionRecord {
key,
owner_client: client,
owner_epoch,
result,
});
disposition
}
#[tokio::test]
async fn concurrent_same_key_start_reuses_one_owner_epoch() {
let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
let starts = Arc::new(AtomicUsize::new(0));
let stopped_owners = Arc::new(Mutex::new(Vec::new()));
let api_key = "pk_test_same_managed_session_owner".to_string();
let (first_entered_tx, first_entered_rx) = oneshot::channel();
let (first_release_tx, first_release_rx) = oneshot::channel();
let first = tokio::spawn(simulate_managed_session_start(
state.clone(),
api_key.clone(),
"same-user".to_string(),
Some((first_entered_tx, first_release_rx)),
starts.clone(),
stopped_owners.clone(),
));
first_entered_rx.await.expect("first start must pause");
let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
let second_state = state.clone();
let second_starts = starts.clone();
let second_stopped_owners = stopped_owners.clone();
let second = tokio::spawn(async move {
second_attempting_tx.send(()).unwrap();
simulate_managed_session_start(
second_state,
api_key,
"same-user".to_string(),
None,
second_starts,
second_stopped_owners,
)
.await
});
second_attempting_rx.await.unwrap();
tokio::task::yield_now().await;
assert_eq!(starts.load(Ordering::SeqCst), 1);
first_release_tx.send(()).unwrap();
assert_eq!(
first.await.expect("first start task"),
ManagedSessionDisposition::Start
);
assert_eq!(
second.await.expect("second start task"),
ManagedSessionDisposition::Reuse
);
assert_eq!(starts.load(Ordering::SeqCst), 1);
let active = state
.managed_session
.lock()
.await
.clone()
.expect("active owner");
assert_eq!(active.owner_epoch, 1);
assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
assert!(stopped_owners.lock().await.is_empty());
}
#[tokio::test]
async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
let starts = Arc::new(AtomicUsize::new(0));
let stopped_owners = Arc::new(Mutex::new(Vec::new()));
let (first_entered_tx, first_entered_rx) = oneshot::channel();
let (first_release_tx, first_release_rx) = oneshot::channel();
let first = tokio::spawn(simulate_managed_session_start(
state.clone(),
"pk_test_owner_a".to_string(),
"first-user".to_string(),
Some((first_entered_tx, first_release_rx)),
starts.clone(),
stopped_owners.clone(),
));
first_entered_rx.await.expect("first start must pause");
let first_owner = state.client();
let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
let second_state = state.clone();
let second_starts = starts.clone();
let second_stopped_owners = stopped_owners.clone();
let second = tokio::spawn(async move {
second_attempting_tx.send(()).unwrap();
simulate_managed_session_start(
second_state,
"pk_test_owner_b".to_string(),
"second-user".to_string(),
None,
second_starts,
second_stopped_owners,
)
.await
});
second_attempting_rx.await.unwrap();
tokio::task::yield_now().await;
assert!(Arc::ptr_eq(&first_owner, &state.client()));
first_release_tx.send(()).unwrap();
assert_eq!(
first.await.expect("first start task"),
ManagedSessionDisposition::Start
);
assert_eq!(
second.await.expect("replacement start task"),
ManagedSessionDisposition::Replace
);
assert_eq!(starts.load(Ordering::SeqCst), 2);
let active = state
.managed_session
.lock()
.await
.clone()
.expect("active owner");
assert_eq!(active.owner_epoch, 2);
assert!(active.key.contains("second-user"));
assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
assert_eq!(
stopped_owners.lock().await.as_slice(),
[Arc::as_ptr(&first_owner) as usize]
);
assert!(!Arc::ptr_eq(&active.owner_client, &first_owner));
}
#[test]
fn derives_app_tag_from_api_key_by_default() {
let config = OpenRtcTauriConfig {
api_key: Some("pk_test_1234567890abcdef".to_string()),
..OpenRtcTauriConfig::default()
};
assert_eq!(config.app_tag(), "app_1234567890abcdef");
}
#[tokio::test]
async fn requested_native_transport_without_installer_keeps_base_route_available() {
let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
let client = state.client();
ensure_requested_native_transports(
&state,
&client,
Some(&requested_ble_config()),
&native_transport_install_context(),
)
.await
.expect("an unavailable optional transport must not fail base Iroh startup");
assert!(state.installed_native_transports.lock().await.is_empty());
}
#[tokio::test]
async fn native_transport_installer_is_idempotent_for_one_client() {
let installs = Arc::new(AtomicUsize::new(0));
let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
.with_native_transport_installer(Arc::new(FakeNativeTransportInstaller {
installs: installs.clone(),
}));
let client = state.client();
let config = requested_ble_config();
ensure_requested_native_transports(
&state,
&client,
Some(&config),
&native_transport_install_context(),
)
.await
.expect("first install");
ensure_requested_native_transports(
&state,
&client,
Some(&config),
&native_transport_install_context(),
)
.await
.expect("idempotent install");
assert_eq!(installs.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn unavailable_native_transport_keeps_base_route_available() {
let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
.with_native_transport_installer(Arc::new(UnavailableNativeTransportInstaller));
let client = state.client();
ensure_requested_native_transports(
&state,
&client,
Some(&requested_ble_config()),
&native_transport_install_context(),
)
.await
.expect("optional transport installation failure must not fail base Iroh startup");
assert!(state.installed_native_transports.lock().await.is_empty());
}
#[test]
fn explicit_app_tag_wins_over_api_key() {
let config = OpenRtcTauriConfig {
api_key: Some("pk_test_1234567890abcdef".to_string()),
app_tag: Some("space::manual".to_string()),
..OpenRtcTauriConfig::default()
};
assert_eq!(config.app_tag(), "space::manual");
}
#[test]
fn derives_space_app_tag_from_space_key() {
let config = OpenRtcTauriConfig {
api_key: Some("pk_test_abc".to_string()),
space_key: Some("room".to_string()),
..OpenRtcTauriConfig::default()
};
assert!(config.app_tag().starts_with("space::"));
assert_eq!(config.app_tag().len(), "space::".len() + 64);
}
#[test]
fn from_env_accepts_vite_project_and_app_tag() {
let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
let config = OpenRtcTauriConfig::from_env();
assert_eq!(config.project_id, "pluto-rtc-prod");
assert_eq!(config.app_tag(), "app_from_vite_env");
match previous_project {
Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
}
match previous_app_tag {
Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
}
}
#[test]
fn native_state_reconfigures_anonymous_client_to_public_space() {
let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
assert_eq!(state.client().app_tag(), "app_anonymous");
let (client, reconfigured) = state.configure_identity(
Some("pk_test_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string()),
Some("assistant-control-plane".to_string()),
None,
);
assert!(reconfigured);
assert!(client.app_tag().starts_with("space::"));
assert_eq!(client.app_tag(), state.client().app_tag());
assert_ne!(client.app_tag(), "app_anonymous");
}
#[test]
fn managed_session_uses_persisted_native_device_id() {
let identity = openrtc::native_device::NativeDeviceIdentity {
device_id: "persisted-native-device".to_string(),
device_name: "Mac".to_string(),
created_at_ms: 1,
updated_at_ms: 1,
name_source: None,
system_info: None,
};
assert_eq!(
managed_session_device_id(Some(" desktop-e2e-native "), &identity),
"persisted-native-device"
);
assert_eq!(
managed_session_device_id(None, &identity),
"persisted-native-device"
);
}
#[test]
fn repeated_managed_session_triggers_converge_on_one_owner() {
let key = "app:user:persisted-native-device";
let cases = [
("react-remount", true, true, true, true),
("hmr", true, true, true, true),
("auth-refresh", true, true, true, true),
("resume", true, true, true, true),
("alias-change", true, true, true, true),
];
for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
assert_eq!(
managed_session_disposition(
Some(key),
true,
active_presence,
active_auto,
key,
requested_presence,
requested_auto,
),
ManagedSessionDisposition::Reuse,
"{label} must reuse the authoritative tuple"
);
}
assert_eq!(
managed_session_disposition(Some(key), true, false, true, key, true, true),
ManagedSessionDisposition::Refresh,
"failed presence startup must retry idempotently"
);
assert_eq!(
managed_session_disposition(
Some(key),
true,
true,
true,
"app:other-user:persisted-native-device",
true,
true,
),
ManagedSessionDisposition::Replace,
"user switching must replace the previous lifecycle owner"
);
assert_eq!(
managed_session_disposition(Some(key), false, true, true, key, true, true),
ManagedSessionDisposition::Replace,
"the same tuple on a different client epoch must replace the previous owner"
);
}
#[test]
fn managed_session_metadata_uses_authoritative_device_id() {
let raw = serde_json::json!({
"deviceId": "persisted-native-device",
"assistantDevice": {
"deviceId": "desktop-e2e-native"
}
})
.to_string();
let metadata =
metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
assert_eq!(parsed["deviceId"], "desktop-e2e-native");
assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
}
}