use std::collections::HashMap;
use std::path::PathBuf;
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";
#[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(&["OPENRTC_APP_TAG"]);
let project_id = first_env(&["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<()>>>,
loop_state: Mutex<Option<PresenceLoopParams>>,
subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
}
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),
loop_state: Mutex::new(None),
subscriptions: Mutex::new(HashMap::new()),
peer_bi_streams: Mutex::new(HashMap::new()),
}
}
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 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 {
requested_local_device_id(requested).unwrap_or_else(|| 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(),
),
}
}
#[derive(Debug, Clone)]
struct PresenceLoopParams {
user_id: String,
device_name: String,
ticket: String,
metadata: Option<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,
}
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();
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_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,
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 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 let Some(params) = state.loop_state.lock().await.clone() {
let client = state.client();
tokio::spawn(async move {
let _ = client
.update_presence(
¶ms.user_id,
¶ms.device_name,
¶ms.ticket,
params.metadata.as_deref(),
)
.await;
});
}
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)
}
#[tauri::command]
async fn search_rtc_devices(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
) -> Result<SearchRtcDevicesResult, String> {
let devices = state
.client()
.search_devices(&user_id)
.await
.map_err(|error| error.to_string())?;
Ok(SearchRtcDevicesResult { devices })
}
#[tauri::command]
async fn search_rtc_devices_with_status(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
) -> Result<SearchRtcDevicesWithStatusResult, String> {
let devices = state
.client()
.devices_with_status(&user_id)
.await
.map_err(|error| error.to_string())?;
Ok(SearchRtcDevicesWithStatusResult { devices })
}
#[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.loop_state.lock().await = Some(PresenceLoopParams {
user_id,
device_name,
ticket,
metadata,
});
Ok(())
}
#[tauri::command]
async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
state.client().stop_presence_loop();
*state.loop_state.lock().await = None;
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 (client, reconfigured) = state.configure_identity(api_key.clone(), space_key.clone(), None);
if reconfigured {
state.replace_connection_state_forwarder(app.clone(), client.clone());
eprintln!(
"[openrtc-tauri] reconfigured native OpenRTC namespace app_tag={}",
client.app_tag()
);
}
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 data_dir = app_data_dir(&app, &state)?;
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={} using logical signaling identity; native runtime identity is {}",
requested, local_device.device_id
);
}
}
let should_start_presence = presence.unwrap_or(true);
let should_start_auto_connect = auto_connect.unwrap_or(true);
let ticket = if should_start_presence || should_start_auto_connect {
Some(
client
.endpoint_ticket_with_token("user-device", 0)
.await
.map_err(|error| error.to_string())?,
)
} else {
None
};
let managed_metadata =
metadata_with_authoritative_device_id(metadata.clone(), &effective_local_device_id);
if should_start_presence {
start_rtc_presence_loop(
state.clone(),
user_id.clone(),
local_node_id.clone(),
local_device.device_name.clone(),
ticket
.clone()
.ok_or_else(|| "managed session expected a ticket".to_string())?,
managed_metadata.clone(),
api_key,
space_key,
)
.await?;
}
if should_start_auto_connect {
client.start_auto_connect(user_id.clone(), effective_local_device_id.clone());
}
let mut logical_local_device = local_device.clone();
logical_local_device.device_id = effective_local_device_id;
Ok(StartRtcManagedSessionResult {
local_node_id,
ticket_scope: ticket.as_ref().map(|_| "user-device".to_string()),
ticket,
presence_started: should_start_presence,
auto_connect_started: should_start_auto_connect,
local_device: logical_local_device,
})
}
#[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();
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,
) -> Result<openrtc::client::ManagedConnectResult, String> {
state
.client()
.connect_device(device_id.as_deref(), &endpoint_ticket)
.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)
}
#[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_transport_only(&peer_id, timeout_ms)
.await?;
send.write_all(&channel_envelope).await?;
Ok((
connection_id,
remote_node_id,
PeerSendStream::plain(send),
PeerRecvStream::plain(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::*;
#[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");
}
#[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 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_prefers_requested_logical_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),
"desktop-e2e-native"
);
assert_eq!(
managed_session_device_id(None, &identity),
"persisted-native-device"
);
}
#[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");
}
}