1use std::collections::HashMap;
2use std::path::PathBuf;
3use std::sync::atomic::{AtomicU64, Ordering};
4use std::sync::{Arc, RwLock};
5
6use futures::StreamExt;
7use openrtc::application_crypto_streams::{PeerRecvStream, PeerSendStream};
8use openrtc::native_node::IncomingStreamType;
9use serde::{Deserialize, Serialize};
10use tauri::{Emitter, Manager, Runtime};
11use tokio::sync::Mutex;
12
13const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
14const RTC_DEVICE_EVENTS: &str = "rtc-device-events";
15const RTC_SESSION_EVENTS: &str = "rtc-session-events";
16const CONNECTION_STATE_CHANGED_EVENT: &str = "connection-state-changed";
17const PEER_BI_STREAM_INCOMING_EVENT: &str = "openrtc://peer-bi-stream/incoming";
18const PEER_BI_STREAM_CHUNK_EVENT: &str = "openrtc://peer-bi-stream/chunk";
19const PEER_BI_STREAM_CLOSED_EVENT: &str = "openrtc://peer-bi-stream/closed";
20const DISCOVERY_COMMAND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
21
22pub type NativeTransportInstallFuture = std::pin::Pin<
23 Box<
24 dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
25 + Send,
26 >,
27>;
28
29#[derive(Debug, Clone)]
30pub struct NativeTransportInstallContext {
31 pub data_dir: PathBuf,
35}
36
37pub trait NativeTransportInstaller: Send + Sync {
41 fn id(&self) -> &'static str;
42 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
43 fn install(
44 &self,
45 client: Arc<openrtc::client::Client>,
46 config: openrtc::client::TransportConfig,
47 context: NativeTransportInstallContext,
48 ) -> NativeTransportInstallFuture;
49}
50
51struct InstalledNativeTransport {
52 client_ptr: usize,
53 _runtime: Box<dyn std::any::Any + Send + Sync>,
54}
55
56#[derive(Debug, Clone, Serialize, Deserialize)]
57#[serde(rename_all = "camelCase")]
58pub struct OpenRtcTauriConfig {
59 #[serde(default = "default_project_id")]
60 pub project_id: String,
61 #[serde(default)]
62 pub api_key: Option<String>,
63 #[serde(default)]
64 pub app_tag: Option<String>,
65 #[serde(default)]
66 pub space_key: Option<String>,
67 #[serde(default)]
68 pub data_dir: Option<PathBuf>,
69 #[serde(default)]
70 pub transport_config: Option<openrtc::client::TransportConfig>,
71}
72
73impl Default for OpenRtcTauriConfig {
74 fn default() -> Self {
75 Self {
76 project_id: default_project_id(),
77 api_key: None,
78 app_tag: None,
79 space_key: None,
80 data_dir: None,
81 transport_config: None,
82 }
83 }
84}
85
86impl OpenRtcTauriConfig {
87 pub fn from_env() -> Self {
88 let api_key = first_env(&[
89 "VITE_OPENRTC_KEY",
90 "VITE_OPENRTC_API_KEY",
91 "VITE_PLUTO_OPENRTC_API_KEY",
92 "OPENRTC_API_KEY",
93 ]);
94 let space_key = first_env(&[
95 "VITE_OPENRTC_SPACE_KEY",
96 "VITE_OPENRTC_SPACE",
97 "OPENRTC_SPACE_KEY",
98 ]);
99 let app_tag = first_env(&[
100 "VITE_PLUTO_OPENRTC_APP_TAG",
101 "VITE_OPENRTC_APP_TAG",
102 "OPENRTC_APP_TAG",
103 ]);
104 let project_id = first_env(&["VITE_OPENRTC_PROJECT_ID", "OPENRTC_PROJECT_ID"])
105 .unwrap_or_else(default_project_id);
106
107 Self {
108 project_id,
109 api_key,
110 app_tag,
111 space_key,
112 data_dir: None,
113 transport_config: None,
114 }
115 }
116
117 pub fn app_tag(&self) -> String {
118 if let Some(app_tag) = self
119 .app_tag
120 .as_deref()
121 .map(str::trim)
122 .filter(|value| !value.is_empty())
123 {
124 return app_tag.to_string();
125 }
126
127 if let (Some(api_key), Some(space_key)) = (
128 self.api_key
129 .as_deref()
130 .map(str::trim)
131 .filter(|value| !value.is_empty()),
132 self.space_key
133 .as_deref()
134 .map(str::trim)
135 .filter(|value| !value.is_empty()),
136 ) {
137 return openrtc::space_app_tag_from_keys(api_key, space_key);
138 }
139
140 self.api_key
141 .as_deref()
142 .map(openrtc::app_tag_from_api_key)
143 .unwrap_or_else(|| "app_anonymous".to_string())
144 }
145}
146
147fn default_project_id() -> String {
148 openrtc::LIVE_PROJECT_ID.to_string()
149}
150
151fn first_env(names: &[&str]) -> Option<String> {
152 names.iter().find_map(|name| {
153 std::env::var(name)
154 .ok()
155 .map(|value| value.trim().to_string())
156 .filter(|value| !value.is_empty())
157 })
158}
159
160#[derive(Default)]
161struct TokenRelayState {
162 auth_token: RwLock<Option<String>>,
163 refresh_token: RwLock<Option<String>>,
164}
165
166impl TokenRelayState {
167 fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
168 let relay = self.clone();
169 Box::new(move || relay.auth_token.read().ok().and_then(|guard| guard.clone()))
170 }
171
172 fn set(&self, auth_token: Option<String>, refresh_token: Option<String>) {
173 if let Ok(mut guard) = self.auth_token.write() {
174 *guard = normalize_token(auth_token);
175 }
176 if let Ok(mut guard) = self.refresh_token.write() {
177 *guard = normalize_token(refresh_token);
178 }
179 }
180}
181
182fn normalize_token(value: Option<String>) -> Option<String> {
183 value
184 .map(|value| value.trim().to_string())
185 .filter(|value| !value.is_empty())
186}
187
188pub struct OpenRtcTauriState {
189 client: RwLock<Arc<openrtc::client::Client>>,
190 config: RwLock<OpenRtcTauriConfig>,
191 token_relay: Arc<TokenRelayState>,
192 data_dir: Option<PathBuf>,
193 connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
194 presence_loop_active: Mutex<bool>,
195 subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
196 peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
197 managed_session_start_guard: Mutex<()>,
198 managed_session_next_owner_epoch: AtomicU64,
199 managed_session: Mutex<Option<ManagedSessionRecord>>,
200 native_transport_installers: Vec<Arc<dyn NativeTransportInstaller>>,
201 installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
202}
203
204impl OpenRtcTauriState {
205 pub fn new(config: OpenRtcTauriConfig) -> Self {
206 openrtc::ensure_default_rustls_provider();
207
208 let token_relay = Arc::new(TokenRelayState::default());
209 let client = build_client(&config, &token_relay);
210 let data_dir = config.data_dir.clone();
211
212 Self {
213 client: RwLock::new(client),
214 config: RwLock::new(config),
215 token_relay,
216 data_dir,
217 connection_state_forwarder: std::sync::Mutex::new(None),
218 presence_loop_active: Mutex::new(false),
219 subscriptions: Mutex::new(HashMap::new()),
220 peer_bi_streams: Mutex::new(HashMap::new()),
221 managed_session_start_guard: Mutex::new(()),
222 managed_session_next_owner_epoch: AtomicU64::new(0),
223 managed_session: Mutex::new(None),
224 native_transport_installers: Vec::new(),
225 installed_native_transports: Mutex::new(HashMap::new()),
226 }
227 }
228
229 pub fn with_native_transport_installer(
230 mut self,
231 installer: Arc<dyn NativeTransportInstaller>,
232 ) -> Self {
233 self.native_transport_installers.push(installer);
234 self
235 }
236
237 pub fn client(&self) -> Arc<openrtc::client::Client> {
238 self.client
239 .read()
240 .map(|guard| guard.clone())
241 .unwrap_or_else(|_| build_client(&OpenRtcTauriConfig::default(), &self.token_relay))
242 }
243
244 pub fn set_auth_token(&self, auth_token: Option<String>, refresh_token: Option<String>) {
245 self.token_relay.set(auth_token, refresh_token);
246 }
247
248 fn configure_identity(
249 &self,
250 api_key: Option<String>,
251 space_key: Option<String>,
252 app_tag: Option<String>,
253 ) -> (Arc<openrtc::client::Client>, bool) {
254 let mut next_config = self
255 .config
256 .read()
257 .map(|guard| guard.clone())
258 .unwrap_or_default();
259 apply_identity_override(&mut next_config.api_key, api_key);
260 apply_identity_override(&mut next_config.space_key, space_key);
261 apply_identity_override(&mut next_config.app_tag, app_tag);
262
263 let next_app_tag = next_config.app_tag();
264 let current = self.client();
265 if current.app_tag() == next_app_tag {
266 if let Ok(mut guard) = self.config.write() {
267 *guard = next_config;
268 }
269 return (current, false);
270 }
271
272 let next_client = build_client(&next_config, &self.token_relay);
273 if let Ok(mut guard) = self.client.write() {
274 *guard = next_client.clone();
275 }
276 if let Ok(mut guard) = self.config.write() {
277 *guard = next_config;
278 }
279 current.stop_auth_scoped_activity();
282 (next_client, true)
283 }
284
285 fn allocate_managed_session_owner_epoch(&self) -> u64 {
286 self.managed_session_next_owner_epoch
287 .fetch_add(1, Ordering::SeqCst)
288 + 1
289 }
290
291 fn replace_connection_state_forwarder<R: Runtime>(
292 &self,
293 app: tauri::AppHandle<R>,
294 client: Arc<openrtc::client::Client>,
295 ) {
296 let next = tauri::async_runtime::spawn(async move {
297 forward_connection_state_events(app, client).await;
298 });
299 if let Ok(mut guard) = self.connection_state_forwarder.lock() {
300 if let Some(previous) = guard.replace(next) {
301 previous.abort();
302 }
303 }
304 }
305}
306
307fn build_client(
308 config: &OpenRtcTauriConfig,
309 token_relay: &Arc<TokenRelayState>,
310) -> Arc<openrtc::client::Client> {
311 let mut builder = openrtc::client::Client::builder_with_app_tag(
312 config.project_id.clone(),
313 config.app_tag(),
314 token_relay.token_provider(),
315 );
316 if let Some(transport_config) = config.transport_config.clone() {
317 builder = builder.transport_config(transport_config);
318 }
319 Arc::new(builder.build())
320}
321
322fn apply_identity_override(target: &mut Option<String>, value: Option<String>) {
323 if let Some(value) = value
324 .as_deref()
325 .map(str::trim)
326 .filter(|value| !value.is_empty())
327 {
328 *target = Some(value.to_string());
329 }
330}
331
332fn requested_local_device_id(value: Option<&str>) -> Option<String> {
333 value
334 .map(str::trim)
335 .filter(|value| !value.is_empty())
336 .map(ToOwned::to_owned)
337}
338
339fn managed_session_device_id(
340 _requested: Option<&str>,
341 native_identity: &openrtc::native_device::NativeDeviceIdentity,
342) -> String {
343 native_identity.device_id.clone()
347}
348
349fn metadata_with_authoritative_device_id(
350 metadata: Option<String>,
351 device_id: &str,
352) -> Option<String> {
353 let device_id = device_id.trim();
354 if device_id.is_empty() {
355 return metadata;
356 }
357
358 let Some(raw_metadata) = metadata else {
359 return Some(serde_json::json!({ "deviceId": device_id }).to_string());
360 };
361
362 match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
363 Ok(serde_json::Value::Object(mut map)) => {
364 map.insert(
365 "deviceId".to_string(),
366 serde_json::Value::String(device_id.to_string()),
367 );
368 Some(serde_json::Value::Object(map).to_string())
369 }
370 _ => Some(
371 serde_json::json!({
372 "deviceId": device_id,
373 "metadata": raw_metadata,
374 })
375 .to_string(),
376 ),
377 }
378}
379
380struct PeerBiStreamHandle {
381 send: Arc<Mutex<Option<PeerSendStream>>>,
382 recv: Option<PeerRecvStream>,
383 read_task: Option<tokio::task::JoinHandle<()>>,
384}
385
386#[derive(Serialize)]
387#[serde(rename_all = "camelCase")]
388pub struct OpenPeerBiStreamResult {
389 stream_id: String,
390 connection_id: Option<String>,
391 remote_node_id: String,
392}
393
394#[derive(Serialize, Clone)]
395#[serde(rename_all = "camelCase")]
396struct IncomingPeerBiStreamEvent {
397 request_id: String,
398 stream_id: String,
399 remote_node_id: String,
400}
401
402#[derive(Serialize, Clone)]
403#[serde(rename_all = "camelCase")]
404struct PeerBiStreamChunkEvent {
405 stream_id: String,
406 bytes: Vec<u8>,
407}
408
409#[derive(Serialize, Clone)]
410#[serde(rename_all = "camelCase")]
411struct PeerBiStreamClosedEvent {
412 stream_id: String,
413 error: Option<String>,
414}
415
416#[derive(Debug, Clone, Serialize)]
417#[serde(rename_all = "camelCase")]
418pub struct SearchRtcDevicesResult {
419 devices: Vec<openrtc::signaling::Device>,
420}
421
422#[derive(Debug, Clone, Serialize)]
423#[serde(rename_all = "camelCase")]
424pub struct SearchRtcDevicesWithStatusResult {
425 devices: Vec<openrtc::client::DeviceStatusSnapshot>,
426}
427
428#[derive(Debug, Clone, Serialize)]
429#[serde(rename_all = "camelCase")]
430pub struct StartRtcManagedSessionResult {
431 local_node_id: String,
432 ticket_scope: Option<String>,
433 ticket: Option<String>,
434 presence_started: bool,
435 auto_connect_started: bool,
436 local_device: openrtc::native_device::NativeDeviceIdentity,
437}
438
439#[derive(Clone)]
440struct ManagedSessionRecord {
441 key: String,
442 owner_client: Arc<openrtc::client::Client>,
443 owner_epoch: u64,
444 result: StartRtcManagedSessionResult,
445}
446
447#[derive(Debug, Clone, Copy, PartialEq, Eq)]
448enum ManagedSessionDisposition {
449 Start,
450 Reuse,
451 Refresh,
452 Replace,
453}
454
455fn managed_session_disposition(
456 active_key: Option<&str>,
457 active_client_matches: bool,
458 active_presence_started: bool,
459 active_auto_connect_started: bool,
460 requested_key: &str,
461 requested_presence: bool,
462 requested_auto_connect: bool,
463) -> ManagedSessionDisposition {
464 let Some(active_key) = active_key else {
465 return ManagedSessionDisposition::Start;
466 };
467 if !active_client_matches || active_key != requested_key {
468 return ManagedSessionDisposition::Replace;
469 }
470 if (!requested_presence || active_presence_started)
471 && (!requested_auto_connect || active_auto_connect_started)
472 {
473 ManagedSessionDisposition::Reuse
474 } else {
475 ManagedSessionDisposition::Refresh
476 }
477}
478
479pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
480 init_with_state(OpenRtcTauriState::new(config))
481}
482
483pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
484 tauri::plugin::Builder::new(PLUGIN_NAME)
485 .setup(move |app, _api| {
486 let client = state.client();
487 let transport_config = state
488 .config
489 .read()
490 .ok()
491 .and_then(|config| config.transport_config.clone());
492 let install_context = NativeTransportInstallContext {
497 data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
498 };
499 tauri::async_runtime::block_on(ensure_requested_native_transports(
500 &state,
501 &client,
502 transport_config.as_ref(),
503 &install_context,
504 ))
505 .map_err(std::io::Error::other)?;
506 state.replace_connection_state_forwarder(app.clone(), client);
507 app.manage(state);
508 Ok(())
509 })
510 .invoke_handler(tauri::generate_handler![
511 desktop_set_pluto_auth_token,
512 openrtc_set_auth_token,
513 rtc_native_status,
514 get_rtc_local_device_info,
515 update_rtc_local_device_name,
516 get_iroh_node_id,
517 start_iroh_node,
518 get_iroh_endpoint_ticket,
519 register_session_token,
520 get_endpoint_ticket_with_token,
521 validate_session_token,
522 revoke_session_tokens_by_scope,
523 search_rtc_devices,
524 search_rtc_devices_with_status,
525 update_rtc_presence,
526 update_rtc_device,
527 delete_rtc_device,
528 set_rtc_offline,
529 start_rtc_presence_loop,
530 stop_rtc_presence_loop,
531 start_rtc_managed_session,
532 start_rtc_auto_connect,
533 stop_rtc_auto_connect,
534 force_rtc_reconnect_snapshot,
535 notify_rtc_network_change,
536 connect_to_device,
537 disconnect_device,
538 set_auto_connect_excluded,
539 resolve_rtc_peer_connection_records,
540 resolve_rtc_peer_identity,
541 get_rtc_peer_session,
542 list_rtc_peer_sessions,
543 list_rtc_managed_connections,
544 wait_for_rtc_settled_peer,
545 list_rtc_connection_states,
546 get_rtc_connection_state,
547 start_rtc_device_subscription,
548 start_rtc_session_subscription,
549 stop_rtc_subscription,
550 start_incoming_peer_bi_streams,
551 open_peer_bi_stream,
552 open_peer_bi_transport_only_stream,
553 open_peer_native_bi_stream,
554 write_peer_bi_stream,
555 start_peer_bi_stream_read,
556 close_peer_bi_stream,
557 get_app_limits,
558 ])
559 .build()
560}
561
562pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
563 init(OpenRtcTauriConfig::from_env())
564}
565
566async fn forward_connection_state_events<R: Runtime>(
567 app: tauri::AppHandle<R>,
568 client: Arc<openrtc::client::Client>,
569) {
570 let mut rx = client.subscribe_native_connection_state_updates();
571 while let Ok(snapshot) = rx.recv().await {
572 let _ = app.emit(CONNECTION_STATE_CHANGED_EVENT, snapshot);
573 }
574}
575
576fn app_data_dir<R: Runtime>(
577 app: &tauri::AppHandle<R>,
578 state: &OpenRtcTauriState,
579) -> Result<PathBuf, String> {
580 state
581 .data_dir
582 .clone()
583 .or_else(|| app.path().app_data_dir().ok())
584 .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
585}
586
587async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
588 if let Some(node_id) = client.current_node_id().await {
589 return Ok(node_id);
590 }
591 client
592 .init_iroh(None, Vec::new())
593 .await
594 .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
595}
596
597async fn ensure_requested_native_transports(
598 state: &OpenRtcTauriState,
599 client: &Arc<openrtc::client::Client>,
600 transports: Option<&openrtc::client::TransportConfig>,
601 context: &NativeTransportInstallContext,
602) -> Result<(), String> {
603 let Some(transports) = transports else {
604 return Ok(());
605 };
606 let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
607 if ble_requested
608 && !state
609 .native_transport_installers
610 .iter()
611 .any(|installer| installer.is_requested(transports))
612 {
613 eprintln!(
614 "[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
615 );
616 return Ok(());
617 }
618 let client_ptr = Arc::as_ptr(client) as usize;
619 for installer in &state.native_transport_installers {
620 if !installer.is_requested(transports) {
621 continue;
622 }
623 let already_installed = state
624 .installed_native_transports
625 .lock()
626 .await
627 .get(installer.id())
628 .is_some_and(|installed| installed.client_ptr == client_ptr);
629 if already_installed {
630 continue;
631 }
632 if client.current_node_id().await.is_some() {
633 eprintln!(
634 "[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
635 installer.id()
636 );
637 continue;
638 }
639 let runtime = match installer
640 .install(client.clone(), transports.clone(), context.clone())
641 .await
642 {
643 Ok(runtime) => runtime,
644 Err(error) => {
645 eprintln!(
646 "[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
647 installer.id(),
648 error
649 );
650 continue;
651 }
652 };
653 state.installed_native_transports.lock().await.insert(
654 installer.id(),
655 InstalledNativeTransport {
656 client_ptr,
657 _runtime: runtime,
658 },
659 );
660 }
661 Ok(())
662}
663
664async fn register_peer_bi_stream<R: Runtime>(
665 app: tauri::AppHandle<R>,
666 state: &OpenRtcTauriState,
667 connection_id: Option<String>,
668 remote_node_id: String,
669 send: PeerSendStream,
670 recv: PeerRecvStream,
671 start_reading: bool,
672) -> OpenPeerBiStreamResult {
673 let stream_id = uuid::Uuid::new_v4().to_string();
674 let send = Arc::new(Mutex::new(Some(send)));
675 let (recv, read_task) = if start_reading {
676 (
677 None,
678 Some(spawn_peer_bi_stream_reader(
679 app.clone(),
680 stream_id.clone(),
681 recv,
682 )),
683 )
684 } else {
685 (Some(recv), None)
689 };
690
691 state.peer_bi_streams.lock().await.insert(
692 stream_id.clone(),
693 PeerBiStreamHandle {
694 send,
695 recv,
696 read_task,
697 },
698 );
699
700 OpenPeerBiStreamResult {
701 stream_id,
702 connection_id,
703 remote_node_id,
704 }
705}
706
707fn spawn_peer_bi_stream_reader<R: Runtime>(
708 app: tauri::AppHandle<R>,
709 stream_id: String,
710 mut recv: PeerRecvStream,
711) -> tokio::task::JoinHandle<()> {
712 tokio::spawn(async move {
713 let mut chunk = vec![0_u8; 64 * 1024];
714 let mut close_error: Option<String> = None;
715 loop {
716 match recv.read(&mut chunk).await {
717 Ok(0) => break,
718 Ok(n) => {
719 let _ = app.emit(
720 PEER_BI_STREAM_CHUNK_EVENT,
721 PeerBiStreamChunkEvent {
722 stream_id: stream_id.clone(),
723 bytes: chunk[..n].to_vec(),
724 },
725 );
726 }
727 Err(error) => {
728 close_error = Some(error.to_string());
729 break;
730 }
731 }
732 }
733 let _ = app.emit(
734 PEER_BI_STREAM_CLOSED_EVENT,
735 PeerBiStreamClosedEvent {
736 stream_id,
737 error: close_error,
738 },
739 );
740 })
741}
742
743async fn open_peer_bi_with<R: Runtime, F, Fut>(
744 app: tauri::AppHandle<R>,
745 state: tauri::State<'_, OpenRtcTauriState>,
746 peer_id: String,
747 timeout_ms: Option<u64>,
748 open: F,
749) -> Result<OpenPeerBiStreamResult, String>
750where
751 F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
752 Fut: std::future::Future<
753 Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
754 >,
755{
756 let peer_id = peer_id.trim().to_string();
757 if peer_id.is_empty() {
758 return Err("peerId is required".to_string());
759 }
760
761 let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
762 .await
763 .map_err(|error| format!("open peer bi stream failed: {error}"))?;
764 Ok(register_peer_bi_stream(app, &state, connection_id, remote_node_id, send, recv, true).await)
765}
766
767#[tauri::command]
768async fn desktop_set_pluto_auth_token(
769 state: tauri::State<'_, OpenRtcTauriState>,
770 auth_token: Option<String>,
771 refresh_token: Option<String>,
772) -> Result<(), String> {
773 state.set_auth_token(auth_token, refresh_token);
774 if *state.presence_loop_active.lock().await {
775 state.client().request_presence_update();
776 }
777 Ok(())
778}
779
780#[tauri::command]
781async fn openrtc_set_auth_token(
782 state: tauri::State<'_, OpenRtcTauriState>,
783 auth_token: Option<String>,
784 refresh_token: Option<String>,
785) -> Result<(), String> {
786 desktop_set_pluto_auth_token(state, auth_token, refresh_token).await
787}
788
789#[tauri::command]
790async fn rtc_native_status(
791 state: tauri::State<'_, OpenRtcTauriState>,
792) -> Result<openrtc::client::RuntimeStatus, String> {
793 Ok(state.client().runtime_status().await)
794}
795
796#[tauri::command]
797async fn get_rtc_local_device_info<R: Runtime>(
798 app: tauri::AppHandle<R>,
799 state: tauri::State<'_, OpenRtcTauriState>,
800) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
801 let data_dir = app_data_dir(&app, &state)?;
802 state
803 .client()
804 .init_native_device_identity(data_dir, None)
805 .await
806 .map_err(|error| error.to_string())
807}
808
809#[tauri::command]
810async fn update_rtc_local_device_name<R: Runtime>(
811 app: tauri::AppHandle<R>,
812 state: tauri::State<'_, OpenRtcTauriState>,
813 device_name: String,
814) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
815 let data_dir = app_data_dir(&app, &state)?;
816 let _ = state
817 .client()
818 .init_native_device_identity(data_dir, None)
819 .await
820 .map_err(|error| error.to_string())?;
821 state
822 .client()
823 .update_native_device_name(&device_name)
824 .await
825 .map_err(|error| error.to_string())
826}
827
828#[tauri::command]
829async fn get_iroh_node_id(
830 state: tauri::State<'_, OpenRtcTauriState>,
831) -> Result<Option<String>, String> {
832 Ok(state.client().current_node_id().await)
833}
834
835#[tauri::command]
836async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
837 ensure_iroh_node(state.client()).await
838}
839
840#[tauri::command]
841async fn get_iroh_endpoint_ticket(
842 state: tauri::State<'_, OpenRtcTauriState>,
843) -> Result<String, String> {
844 ensure_iroh_node(state.client()).await?;
845 state
846 .client()
847 .endpoint_ticket()
848 .await
849 .map_err(|error| error.to_string())
850}
851
852#[tauri::command]
853async fn register_session_token(
854 state: tauri::State<'_, OpenRtcTauriState>,
855 token: String,
856 scope: String,
857 max_connections: u32,
858 expires_at_ms: Option<u64>,
859) -> Result<(), String> {
860 if let Some(expires_at_ms) = expires_at_ms {
861 state.client().register_session_token_with_expiry_ms(
862 token,
863 scope,
864 max_connections,
865 expires_at_ms,
866 );
867 } else {
868 state
869 .client()
870 .register_session_token(token, scope, max_connections);
871 }
872 Ok(())
873}
874
875#[tauri::command]
876async fn get_endpoint_ticket_with_token(
877 state: tauri::State<'_, OpenRtcTauriState>,
878 scope: String,
879 max_connections: u32,
880) -> Result<String, String> {
881 ensure_iroh_node(state.client()).await?;
882 state
883 .client()
884 .endpoint_ticket_with_token(&scope, max_connections)
885 .await
886 .map_err(|error| error.to_string())
887}
888
889#[tauri::command]
890async fn validate_session_token(
891 state: tauri::State<'_, OpenRtcTauriState>,
892 token: String,
893 connection_id: Option<String>,
894) -> Result<String, String> {
895 if let Some(connection_id) = connection_id
896 .as_deref()
897 .map(str::trim)
898 .filter(|value| !value.is_empty())
899 {
900 state
901 .client()
902 .validate_session_token_for_connection(&token, connection_id)
903 .await
904 .map_err(|error| error.to_string())
905 } else {
906 state.client().validate_session_token(&token)
907 }
908}
909
910#[tauri::command]
911async fn revoke_session_tokens_by_scope(
912 state: tauri::State<'_, OpenRtcTauriState>,
913 scope: String,
914) -> Result<Vec<String>, String> {
915 Ok(state.client().revoke_tokens_by_scope(&scope).await)
916}
917
918fn discovery_verbose() -> bool {
919 std::env::var("OPENRTC_SIGNALING_VERBOSE")
920 .map(|value| {
921 matches!(
922 value.trim().to_ascii_lowercase().as_str(),
923 "1" | "true" | "yes" | "on"
924 )
925 })
926 .unwrap_or(false)
927}
928
929#[tauri::command]
930async fn search_rtc_devices(
931 state: tauri::State<'_, OpenRtcTauriState>,
932 user_id: String,
933) -> Result<SearchRtcDevicesResult, String> {
934 let client = state.client();
935 let app_tag = client.app_tag().to_string();
936 let started = std::time::Instant::now();
937 if discovery_verbose() {
938 eprintln!(
939 "[openrtc-tauri][discovery] search_rtc_devices start user_id={} app_tag={}",
940 user_id, app_tag
941 );
942 }
943 let devices = match tokio::time::timeout(
944 DISCOVERY_COMMAND_TIMEOUT,
945 client.search_devices(&user_id),
946 )
947 .await
948 {
949 Ok(Ok(devices)) => devices,
950 Ok(Err(error)) => {
951 eprintln!(
952 "[openrtc-tauri][discovery] search_rtc_devices failed user_id={} app_tag={} elapsed_ms={} error={}",
953 user_id,
954 app_tag,
955 started.elapsed().as_millis(),
956 error
957 );
958 return Err(error.to_string());
959 }
960 Err(_) => {
961 let message = format!(
962 "search_rtc_devices native discovery timed out after {}ms user_id={} app_tag={}",
963 DISCOVERY_COMMAND_TIMEOUT.as_millis(),
964 user_id,
965 app_tag
966 );
967 eprintln!(
968 "[openrtc-tauri][discovery] {} elapsed_ms={}",
969 message,
970 started.elapsed().as_millis()
971 );
972 return Err(message);
973 }
974 };
975 if discovery_verbose() {
976 eprintln!(
977 "[openrtc-tauri][discovery] search_rtc_devices done user_id={} app_tag={} count={} elapsed_ms={}",
978 user_id,
979 app_tag,
980 devices.len(),
981 started.elapsed().as_millis()
982 );
983 }
984 Ok(SearchRtcDevicesResult { devices })
985}
986
987#[tauri::command]
988async fn search_rtc_devices_with_status(
989 state: tauri::State<'_, OpenRtcTauriState>,
990 user_id: String,
991) -> Result<SearchRtcDevicesWithStatusResult, String> {
992 let client = state.client();
993 let app_tag = client.app_tag().to_string();
994 let started = std::time::Instant::now();
995 if discovery_verbose() {
996 eprintln!(
997 "[openrtc-tauri][discovery] search_rtc_devices_with_status start user_id={} app_tag={}",
998 user_id, app_tag
999 );
1000 }
1001 let devices = match tokio::time::timeout(
1002 DISCOVERY_COMMAND_TIMEOUT,
1003 client.devices_with_status(&user_id),
1004 )
1005 .await
1006 {
1007 Ok(Ok(devices)) => devices,
1008 Ok(Err(error)) => {
1009 eprintln!(
1010 "[openrtc-tauri][discovery] search_rtc_devices_with_status failed user_id={} app_tag={} elapsed_ms={} error={}",
1011 user_id,
1012 app_tag,
1013 started.elapsed().as_millis(),
1014 error
1015 );
1016 return Err(error.to_string());
1017 }
1018 Err(_) => {
1019 let message = format!(
1020 "search_rtc_devices_with_status native discovery timed out after {}ms user_id={} app_tag={}",
1021 DISCOVERY_COMMAND_TIMEOUT.as_millis(),
1022 user_id,
1023 app_tag
1024 );
1025 eprintln!(
1026 "[openrtc-tauri][discovery] {} elapsed_ms={}",
1027 message,
1028 started.elapsed().as_millis()
1029 );
1030 return Err(message);
1031 }
1032 };
1033 if discovery_verbose() {
1034 eprintln!(
1035 "[openrtc-tauri][discovery] search_rtc_devices_with_status done user_id={} app_tag={} count={} elapsed_ms={}",
1036 user_id,
1037 app_tag,
1038 devices.len(),
1039 started.elapsed().as_millis()
1040 );
1041 }
1042 Ok(SearchRtcDevicesWithStatusResult { devices })
1043}
1044
1045#[tauri::command]
1046async fn update_rtc_presence(
1047 state: tauri::State<'_, OpenRtcTauriState>,
1048 user_id: String,
1049 device_name: String,
1050 ticket: String,
1051 metadata: Option<String>,
1052) -> Result<(), String> {
1053 state
1054 .client()
1055 .update_durable_device_record_with_ttl(
1056 &user_id,
1057 &device_name,
1058 &ticket,
1059 7 * 24 * 60 * 60 * 1000,
1060 metadata.as_deref(),
1061 )
1062 .await
1063 .map_err(|error| error.to_string())?;
1064 state
1065 .client()
1066 .update_live_presence_record(&user_id, &device_name, &ticket, metadata.as_deref())
1067 .await
1068 .map_err(|error| error.to_string())
1069}
1070
1071#[tauri::command]
1072async fn update_rtc_device(
1073 state: tauri::State<'_, OpenRtcTauriState>,
1074 user_id: String,
1075 device_id: String,
1076 device_name: Option<String>,
1077 capabilities: Option<openrtc::signaling::DeviceCapabilities>,
1078 metadata: Option<String>,
1079) -> Result<(), String> {
1080 state
1081 .client()
1082 .update_device(
1083 &user_id,
1084 &device_id,
1085 device_name.as_deref(),
1086 capabilities,
1087 metadata.as_deref(),
1088 )
1089 .await
1090 .map_err(|error| error.to_string())
1091}
1092
1093#[tauri::command]
1094async fn delete_rtc_device(
1095 state: tauri::State<'_, OpenRtcTauriState>,
1096 user_id: String,
1097 device_id: String,
1098) -> Result<(), String> {
1099 state
1100 .client()
1101 .delete_device(&user_id, &device_id)
1102 .await
1103 .map_err(|error| error.to_string())
1104}
1105
1106#[tauri::command]
1107async fn set_rtc_offline(
1108 state: tauri::State<'_, OpenRtcTauriState>,
1109 user_id: String,
1110) -> Result<(), String> {
1111 state
1112 .client()
1113 .set_offline(&user_id)
1114 .await
1115 .map_err(|error| error.to_string())
1116}
1117
1118#[tauri::command]
1119async fn start_rtc_presence_loop(
1120 state: tauri::State<'_, OpenRtcTauriState>,
1121 user_id: String,
1122 local_node_id: String,
1123 device_name: String,
1124 ticket: String,
1125 metadata: Option<String>,
1126 api_key: Option<String>,
1127 space_key: Option<String>,
1128) -> Result<(), String> {
1129 let _ = local_node_id;
1130 let (client, _) = state.configure_identity(api_key.clone(), space_key, None);
1131 if let Some(api_key) = api_key
1132 .as_deref()
1133 .map(str::trim)
1134 .filter(|value| !value.is_empty())
1135 {
1136 let api_key = api_key.to_string();
1137 let client = client.clone();
1138 tokio::spawn(async move {
1139 let _ = client.fetch_and_cache_app_limits(&api_key).await;
1140 });
1141 }
1142 client.start_signaling_loop(
1143 user_id.clone(),
1144 device_name.clone(),
1145 ticket.clone(),
1146 metadata.clone(),
1147 );
1148 *state.presence_loop_active.lock().await = true;
1149 Ok(())
1150}
1151
1152#[tauri::command]
1153async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1154 state.client().stop_presence_loop();
1155 *state.presence_loop_active.lock().await = false;
1156 if let Some(active) = state.managed_session.lock().await.as_mut() {
1157 active.result.presence_started = false;
1158 }
1159 stop_subscription_by_id(&state, "internal-presence-loop").await;
1160 Ok(())
1161}
1162
1163#[tauri::command]
1164async fn start_rtc_managed_session<R: Runtime>(
1165 app: tauri::AppHandle<R>,
1166 state: tauri::State<'_, OpenRtcTauriState>,
1167 user_id: String,
1168 device_name: Option<String>,
1169 local_device_id: Option<String>,
1170 metadata: Option<String>,
1171 transports: Option<openrtc::client::TransportConfig>,
1172 auto_connect: Option<bool>,
1173 presence: Option<bool>,
1174 api_key: Option<String>,
1175 space_key: Option<String>,
1176) -> Result<StartRtcManagedSessionResult, String> {
1177 let command_started = std::time::Instant::now();
1178 let _start_guard = state.managed_session_start_guard.lock().await;
1179 let (client, reconfigured) = state.configure_identity(api_key.clone(), space_key.clone(), None);
1180 let app_tag = client.app_tag().to_string();
1181 eprintln!(
1182 "[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
1183 user_id,
1184 app_tag,
1185 device_name
1186 .as_deref()
1187 .map(str::trim)
1188 .map(|value| !value.is_empty())
1189 .unwrap_or(false),
1190 local_device_id
1191 .as_deref()
1192 .map(str::trim)
1193 .map(|value| !value.is_empty())
1194 .unwrap_or(false),
1195 metadata
1196 .as_deref()
1197 .map(str::trim)
1198 .map(|value| !value.is_empty())
1199 .unwrap_or(false),
1200 auto_connect.unwrap_or(true),
1201 presence.unwrap_or(true)
1202 );
1203 if reconfigured {
1204 state.replace_connection_state_forwarder(app.clone(), client.clone());
1205 eprintln!(
1206 "[openrtc-tauri] reconfigured native OpenRTC namespace app_tag={}",
1207 client.app_tag()
1208 );
1209 }
1210
1211 let data_dir = app_data_dir(&app, &state)?;
1212 let install_context = NativeTransportInstallContext {
1213 data_dir: data_dir.clone(),
1214 };
1215 ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
1216 .await?;
1217 if let Some(transport_config) = transports {
1218 client
1219 .update_transport_config(transport_config)
1220 .await
1221 .map_err(|error| error.to_string())?;
1222 }
1223
1224 let local_node_id = ensure_iroh_node(client.clone()).await?;
1225 let local_device = client
1226 .init_native_device_identity(data_dir, device_name.as_deref())
1227 .await
1228 .map_err(|error| error.to_string())?;
1229
1230 let effective_local_device_id =
1231 managed_session_device_id(local_device_id.as_deref(), &local_device);
1232 if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
1233 if requested != local_device.device_id {
1234 eprintln!(
1235 "[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
1236 requested, local_device.device_id
1237 );
1238 }
1239 }
1240
1241 let should_start_presence = presence.unwrap_or(true);
1242 let should_start_auto_connect = auto_connect.unwrap_or(true);
1243 let managed_session_key = format!("{}:{}:{}", app_tag, user_id, effective_local_device_id);
1244 let (disposition, reused) = {
1245 let active = state.managed_session.lock().await;
1246 let disposition = managed_session_disposition(
1247 active.as_ref().map(|record| record.key.as_str()),
1248 active
1249 .as_ref()
1250 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1251 active
1252 .as_ref()
1253 .is_some_and(|record| record.result.presence_started),
1254 active
1255 .as_ref()
1256 .is_some_and(|record| record.result.auto_connect_started),
1257 &managed_session_key,
1258 should_start_presence,
1259 should_start_auto_connect,
1260 );
1261 let reused = (disposition == ManagedSessionDisposition::Reuse).then(|| {
1262 active
1263 .as_ref()
1264 .expect("managed session reuse requires an active record")
1265 .result
1266 .clone()
1267 });
1268 (disposition, reused)
1269 };
1270 if let Some(reused) = reused {
1271 eprintln!(
1272 "[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
1273 managed_session_key,
1274 command_started.elapsed().as_millis()
1275 );
1276 return Ok(reused);
1277 }
1278
1279 if disposition == ManagedSessionDisposition::Replace {
1282 let previous = state
1283 .managed_session
1284 .lock()
1285 .await
1286 .take()
1287 .expect("managed session replacement requires an active record");
1288 previous.owner_client.stop_presence_loop();
1289 previous.owner_client.stop_auto_connect();
1290 *state.presence_loop_active.lock().await = false;
1291 }
1292 let ticket_scope = should_start_presence.then(|| "user-device".to_string());
1293 let managed_metadata =
1294 metadata_with_authoritative_device_id(metadata.clone(), &effective_local_device_id);
1295
1296 if should_start_presence {
1297 client
1298 .clone()
1299 .start_managed_user_device_presence_loop(
1300 user_id.clone(),
1301 local_device.device_name.clone(),
1302 managed_metadata.clone(),
1303 )
1304 .await
1305 .map_err(|error| {
1306 format!(
1307 "managed native presence failed initial publication: {}",
1308 error
1309 )
1310 })?;
1311 *state.presence_loop_active.lock().await = true;
1312 if let Some(api_key) = api_key
1313 .as_deref()
1314 .map(str::trim)
1315 .filter(|value| !value.is_empty())
1316 {
1317 let api_key = api_key.to_string();
1318 let limits_client = client.clone();
1319 tokio::spawn(async move {
1320 let _ = limits_client.fetch_and_cache_app_limits(&api_key).await;
1321 });
1322 }
1323 }
1324
1325 if should_start_auto_connect {
1326 client
1327 .clone()
1328 .start_auto_connect(user_id.clone(), effective_local_device_id.clone());
1329 }
1330
1331 let logical_local_device = local_device.clone();
1332
1333 let owner_epoch = match disposition {
1334 ManagedSessionDisposition::Refresh => {
1335 state
1336 .managed_session
1337 .lock()
1338 .await
1339 .as_ref()
1340 .expect("managed session refresh requires an active record")
1341 .owner_epoch
1342 }
1343 ManagedSessionDisposition::Start | ManagedSessionDisposition::Replace => {
1344 state.allocate_managed_session_owner_epoch()
1345 }
1346 ManagedSessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
1347 };
1348
1349 eprintln!(
1350 "[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={}",
1351 user_id,
1352 app_tag,
1353 local_node_id,
1354 logical_local_device.device_id,
1355 owner_epoch,
1356 should_start_presence,
1357 should_start_auto_connect,
1358 ticket_scope.as_deref().unwrap_or("unrestricted"),
1359 command_started.elapsed().as_millis()
1360 );
1361
1362 let result = StartRtcManagedSessionResult {
1363 local_node_id,
1364 ticket_scope,
1365 ticket: None,
1366 presence_started: should_start_presence,
1367 auto_connect_started: should_start_auto_connect,
1368 local_device: logical_local_device,
1369 };
1370 *state.managed_session.lock().await = Some(ManagedSessionRecord {
1371 key: managed_session_key,
1372 owner_client: client,
1373 owner_epoch,
1374 result: result.clone(),
1375 });
1376 Ok(result)
1377}
1378
1379#[tauri::command]
1380async fn start_rtc_auto_connect(
1381 state: tauri::State<'_, OpenRtcTauriState>,
1382 user_id: String,
1383 local_device_id: String,
1384) -> Result<(), String> {
1385 state
1386 .client()
1387 .clone()
1388 .start_auto_connect(user_id, local_device_id);
1389 Ok(())
1390}
1391
1392#[tauri::command]
1393async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1394 state.client().stop_auto_connect();
1395 if let Some(active) = state.managed_session.lock().await.as_mut() {
1396 active.result.auto_connect_started = false;
1397 }
1398 Ok(())
1399}
1400
1401#[tauri::command]
1402async fn force_rtc_reconnect_snapshot(
1403 state: tauri::State<'_, OpenRtcTauriState>,
1404) -> Result<(), String> {
1405 state.client().clone().force_reconnect_snapshot();
1406 Ok(())
1407}
1408
1409#[tauri::command]
1410async fn connect_to_device(
1411 state: tauri::State<'_, OpenRtcTauriState>,
1412 device_id: Option<String>,
1413 endpoint_ticket: String,
1414 timeout_ms: Option<u64>,
1415) -> Result<openrtc::client::ManagedConnectResult, String> {
1416 let client = state.client();
1417 let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
1418 match timeout_ms {
1419 Some(timeout_ms) => {
1420 tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
1421 .await
1422 .map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
1423 .map_err(|error| error.to_string())
1424 }
1425 None => connect.await.map_err(|error| error.to_string()),
1426 }
1427}
1428
1429#[tauri::command]
1430async fn disconnect_device(
1431 state: tauri::State<'_, OpenRtcTauriState>,
1432 device_id: String,
1433 node_id_hint: Option<String>,
1434) -> Result<(), String> {
1435 state
1436 .client()
1437 .disconnect_device(&device_id, node_id_hint.as_deref())
1438 .await;
1439 Ok(())
1440}
1441
1442#[tauri::command]
1443async fn set_auto_connect_excluded(
1444 state: tauri::State<'_, OpenRtcTauriState>,
1445 device_id: String,
1446 excluded: bool,
1447) -> Result<(), String> {
1448 if excluded {
1449 state.client().exclude_peer_and_publish(&device_id).await;
1450 } else {
1451 state.client().unexclude_peer_and_publish(&device_id).await;
1452 }
1453 Ok(())
1454}
1455
1456#[tauri::command]
1457async fn resolve_rtc_peer_connection_records(
1458 state: tauri::State<'_, OpenRtcTauriState>,
1459 id: String,
1460) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1461 Ok(state.client().resolve_peer_connection_records(&id).await)
1462}
1463
1464#[tauri::command]
1465async fn resolve_rtc_peer_identity(
1466 state: tauri::State<'_, OpenRtcTauriState>,
1467 id: String,
1468) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
1469 Ok(state.client().peer_snapshot(&id).await)
1470}
1471
1472#[tauri::command]
1473async fn get_rtc_peer_session(
1474 state: tauri::State<'_, OpenRtcTauriState>,
1475 id: String,
1476) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1477 Ok(state.client().peer_session(&id).await)
1478}
1479
1480#[tauri::command]
1481async fn list_rtc_peer_sessions(
1482 state: tauri::State<'_, OpenRtcTauriState>,
1483) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
1484 Ok(state.client().peer_sessions().await)
1485}
1486
1487#[tauri::command]
1488async fn list_rtc_managed_connections(
1489 state: tauri::State<'_, OpenRtcTauriState>,
1490) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1491 Ok(state.client().list_managed_connections().await)
1492}
1493
1494#[derive(Debug, Clone, Serialize)]
1495#[serde(rename_all = "camelCase")]
1496struct NotifyRtcNetworkChangeResult {
1497 retired_stale_connections: usize,
1498}
1499
1500#[tauri::command]
1501async fn notify_rtc_network_change(
1502 state: tauri::State<'_, OpenRtcTauriState>,
1503) -> Result<NotifyRtcNetworkChangeResult, String> {
1504 let retired_stale_connections = state
1505 .client()
1506 .notify_network_change()
1507 .await
1508 .map_err(|error| error.to_string())?;
1509 Ok(NotifyRtcNetworkChangeResult {
1510 retired_stale_connections,
1511 })
1512}
1513
1514#[tauri::command]
1515async fn wait_for_rtc_settled_peer(
1516 state: tauri::State<'_, OpenRtcTauriState>,
1517 id: String,
1518 timeout_ms: Option<u64>,
1519) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1520 Ok(state.client().wait_for_settled_peer(&id, timeout_ms).await)
1521}
1522
1523#[tauri::command]
1524async fn list_rtc_connection_states(
1525 state: tauri::State<'_, OpenRtcTauriState>,
1526) -> Result<Vec<openrtc::client::ConnectionStateSnapshot>, String> {
1527 Ok(state.client().connection_states().await)
1528}
1529
1530#[tauri::command]
1531async fn get_rtc_connection_state(
1532 state: tauri::State<'_, OpenRtcTauriState>,
1533 connection_id: String,
1534) -> Result<Option<openrtc::client::ConnectionStateSnapshot>, String> {
1535 Ok(state.client().connection_state(&connection_id).await)
1536}
1537
1538#[tauri::command]
1539async fn start_rtc_device_subscription<R: Runtime>(
1540 state: tauri::State<'_, OpenRtcTauriState>,
1541 app: tauri::AppHandle<R>,
1542 request_id: String,
1543 user_id: String,
1544) -> Result<(), String> {
1545 stop_subscription_by_id(&state, &request_id).await;
1546 let mut stream = state
1547 .client()
1548 .subscribe_devices(&user_id)
1549 .await
1550 .map_err(|error| error.to_string())?;
1551 let request_for_task = request_id.clone();
1552 let handle = tokio::spawn(async move {
1553 while let Some(result) = stream.next().await {
1554 let payload = match result {
1555 Ok(events) => serde_json::json!({
1556 "requestId": request_for_task,
1557 "type": "deviceEvents",
1558 "events": events,
1559 }),
1560 Err(error) => serde_json::json!({
1561 "requestId": request_for_task,
1562 "type": "error",
1563 "message": error.to_string(),
1564 }),
1565 };
1566 let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
1567 let _ = app.emit(RTC_DEVICE_EVENTS, payload);
1568 if is_error {
1569 break;
1570 }
1571 }
1572 });
1573 state.subscriptions.lock().await.insert(request_id, handle);
1574 Ok(())
1575}
1576
1577#[tauri::command]
1578async fn start_rtc_session_subscription<R: Runtime>(
1579 state: tauri::State<'_, OpenRtcTauriState>,
1580 app: tauri::AppHandle<R>,
1581 request_id: String,
1582 local_device_id: String,
1583) -> Result<(), String> {
1584 stop_subscription_by_id(&state, &request_id).await;
1585 let mut stream = state
1586 .client()
1587 .subscribe_sessions(&local_device_id)
1588 .await
1589 .map_err(|error| error.to_string())?;
1590 let request_for_task = request_id.clone();
1591 let handle = tokio::spawn(async move {
1592 while let Some(result) = stream.next().await {
1593 let payload = match result {
1594 Ok(events) => serde_json::json!({
1595 "requestId": request_for_task,
1596 "type": "sessionEvents",
1597 "events": events,
1598 }),
1599 Err(error) => serde_json::json!({
1600 "requestId": request_for_task,
1601 "type": "error",
1602 "message": error.to_string(),
1603 }),
1604 };
1605 let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
1606 let _ = app.emit(RTC_SESSION_EVENTS, payload);
1607 if is_error {
1608 break;
1609 }
1610 }
1611 });
1612 state.subscriptions.lock().await.insert(request_id, handle);
1613 Ok(())
1614}
1615
1616#[tauri::command]
1617async fn stop_rtc_subscription(
1618 state: tauri::State<'_, OpenRtcTauriState>,
1619 request_id: String,
1620) -> Result<(), String> {
1621 stop_subscription_by_id(&state, &request_id).await;
1622 Ok(())
1623}
1624
1625async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
1626 if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
1627 handle.abort();
1628 }
1629}
1630
1631#[tauri::command]
1632async fn start_incoming_peer_bi_streams<R: Runtime>(
1633 app: tauri::AppHandle<R>,
1634 state: tauri::State<'_, OpenRtcTauriState>,
1635 request_id: Option<String>,
1636) -> Result<String, String> {
1637 let request_id = request_id
1638 .map(|value| value.trim().to_string())
1639 .filter(|value| !value.is_empty())
1640 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1641
1642 stop_subscription_by_id(&state, &request_id).await;
1643
1644 let incoming = state
1645 .client()
1646 .incoming_streams()
1647 .await
1648 .map_err(|error| format!("failed to subscribe to incoming peer streams: {error}"))?;
1649 let app_for_task = app.clone();
1650 let request_for_task = request_id.clone();
1651 let handle = tokio::spawn(async move {
1652 while let Ok(incoming_stream) = incoming.recv().await {
1653 let remote_node_id = incoming_stream.endpoint_id.to_string();
1654 let IncomingStreamType::Bi(send, recv) = incoming_stream.stream else {
1655 continue;
1656 };
1657 let state_for_task = app_for_task.state::<OpenRtcTauriState>();
1658 let result = register_peer_bi_stream(
1659 app_for_task.clone(),
1660 state_for_task.inner(),
1661 None,
1662 remote_node_id.clone(),
1663 PeerSendStream::plain(send),
1664 PeerRecvStream::plain(recv),
1665 false,
1666 )
1667 .await;
1668 let _ = app_for_task.emit(
1669 PEER_BI_STREAM_INCOMING_EVENT,
1670 IncomingPeerBiStreamEvent {
1671 request_id: request_for_task.clone(),
1672 stream_id: result.stream_id,
1673 remote_node_id,
1674 },
1675 );
1676 }
1677 });
1678
1679 state
1680 .subscriptions
1681 .lock()
1682 .await
1683 .insert(request_id.clone(), handle);
1684 Ok(request_id)
1685}
1686
1687#[tauri::command]
1688async fn open_peer_bi_stream<R: Runtime>(
1689 app: tauri::AppHandle<R>,
1690 state: tauri::State<'_, OpenRtcTauriState>,
1691 peer_id: String,
1692 timeout_ms: Option<u64>,
1693) -> Result<OpenPeerBiStreamResult, String> {
1694 open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
1695 client.open_peer_bi(&peer_id, timeout_ms).await
1696 })
1697 .await
1698}
1699
1700#[tauri::command]
1701async fn open_peer_bi_transport_only_stream<R: Runtime>(
1702 app: tauri::AppHandle<R>,
1703 state: tauri::State<'_, OpenRtcTauriState>,
1704 peer_id: String,
1705 timeout_ms: Option<u64>,
1706) -> Result<OpenPeerBiStreamResult, String> {
1707 open_peer_bi_with(
1708 app,
1709 state,
1710 peer_id,
1711 timeout_ms,
1712 |client, peer_id, timeout_ms| async move {
1713 client
1714 .open_peer_bi_transport_only(&peer_id, timeout_ms)
1715 .await
1716 .map(|(connection_id, remote_node_id, send, recv)| {
1717 (
1718 connection_id,
1719 remote_node_id,
1720 PeerSendStream::plain(send),
1721 PeerRecvStream::plain(recv),
1722 )
1723 })
1724 },
1725 )
1726 .await
1727}
1728
1729#[tauri::command]
1730async fn open_peer_native_bi_stream<R: Runtime>(
1731 app: tauri::AppHandle<R>,
1732 state: tauri::State<'_, OpenRtcTauriState>,
1733 peer_id: String,
1734 label: String,
1735 timeout_ms: Option<u64>,
1736) -> Result<OpenPeerBiStreamResult, String> {
1737 let channel_envelope = openrtc::stream_metadata::encode_channel_envelope(&label, None)
1738 .map_err(|error| error.to_string())?;
1739 open_peer_bi_with(
1740 app,
1741 state,
1742 peer_id,
1743 timeout_ms,
1744 move |client, peer_id, timeout_ms| async move {
1745 let (connection_id, remote_node_id, mut send, recv) =
1746 client.open_peer_bi(&peer_id, timeout_ms).await?;
1747 send.write_all(&channel_envelope).await?;
1748 Ok((connection_id, remote_node_id, send, recv))
1749 },
1750 )
1751 .await
1752}
1753
1754#[tauri::command]
1755async fn write_peer_bi_stream(
1756 state: tauri::State<'_, OpenRtcTauriState>,
1757 stream_id: String,
1758 bytes: Vec<u8>,
1759) -> Result<(), String> {
1760 let stream_id = stream_id.trim().to_string();
1761 if stream_id.is_empty() {
1762 return Err("streamId is required".to_string());
1763 }
1764
1765 let send = {
1766 let streams = state.peer_bi_streams.lock().await;
1767 streams
1768 .get(&stream_id)
1769 .map(|handle| handle.send.clone())
1770 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
1771 };
1772 let result = send
1773 .lock()
1774 .await
1775 .as_mut()
1776 .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
1777 .write_all(&bytes)
1778 .await
1779 .map_err(|error| format!("write peer bi stream failed: {error}"));
1780 result
1781}
1782
1783#[tauri::command]
1784async fn start_peer_bi_stream_read<R: Runtime>(
1785 app: tauri::AppHandle<R>,
1786 state: tauri::State<'_, OpenRtcTauriState>,
1787 stream_id: String,
1788) -> Result<(), String> {
1789 let stream_id = stream_id.trim().to_string();
1790 if stream_id.is_empty() {
1791 return Err("streamId is required".to_string());
1792 }
1793
1794 let mut streams = state.peer_bi_streams.lock().await;
1795 let handle = streams
1796 .get_mut(&stream_id)
1797 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
1798 if handle.read_task.is_some() {
1799 return Ok(());
1800 }
1801 let recv = handle
1802 .recv
1803 .take()
1804 .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
1805 handle.read_task = Some(spawn_peer_bi_stream_reader(app, stream_id, recv));
1806 Ok(())
1807}
1808
1809#[tauri::command]
1810async fn close_peer_bi_stream(
1811 state: tauri::State<'_, OpenRtcTauriState>,
1812 stream_id: String,
1813) -> Result<(), String> {
1814 let stream_id = stream_id.trim().to_string();
1815 if stream_id.is_empty() {
1816 return Err("streamId is required".to_string());
1817 }
1818
1819 let handle = {
1820 let mut streams = state.peer_bi_streams.lock().await;
1821 streams.remove(&stream_id)
1822 };
1823 if let Some(handle) = handle {
1824 let mut send = handle.send.lock().await;
1827 if let Some(send) = send.take() {
1828 let _ = send.finish();
1829 }
1830 drop(send);
1831 if let Some(read_task) = handle.read_task {
1832 read_task.abort();
1833 }
1834 }
1835 Ok(())
1836}
1837
1838#[tauri::command]
1839async fn get_app_limits() -> Result<serde_json::Value, String> {
1840 Ok(serde_json::json!({
1841 "devicesPerUser": -1,
1842 "maxRooms": -1,
1843 "maxMembersPerRoom": -1
1844 }))
1845}
1846
1847#[cfg(test)]
1848mod tests {
1849 use super::*;
1850 use std::sync::atomic::{AtomicUsize, Ordering};
1851 use tokio::sync::oneshot;
1852
1853 struct FakeNativeTransportInstaller {
1854 installs: Arc<AtomicUsize>,
1855 }
1856
1857 impl NativeTransportInstaller for FakeNativeTransportInstaller {
1858 fn id(&self) -> &'static str {
1859 "ble"
1860 }
1861
1862 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
1863 config.ble.as_ref().is_some_and(|ble| ble.enabled)
1864 }
1865
1866 fn install(
1867 &self,
1868 _client: Arc<openrtc::client::Client>,
1869 _config: openrtc::client::TransportConfig,
1870 _context: NativeTransportInstallContext,
1871 ) -> NativeTransportInstallFuture {
1872 self.installs.fetch_add(1, Ordering::SeqCst);
1873 Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
1874 }
1875 }
1876
1877 #[derive(Debug)]
1878 struct UnavailableNativeTransportInstaller;
1879
1880 impl NativeTransportInstaller for UnavailableNativeTransportInstaller {
1881 fn id(&self) -> &'static str {
1882 "ble"
1883 }
1884
1885 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
1886 config.ble.as_ref().is_some_and(|ble| ble.enabled)
1887 }
1888
1889 fn install(
1890 &self,
1891 _client: Arc<openrtc::client::Client>,
1892 _config: openrtc::client::TransportConfig,
1893 _context: NativeTransportInstallContext,
1894 ) -> NativeTransportInstallFuture {
1895 Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
1896 }
1897 }
1898
1899 fn requested_ble_config() -> openrtc::client::TransportConfig {
1900 openrtc::client::TransportConfig {
1901 ble: Some(openrtc::client::BleConfig {
1902 enabled: true,
1903 ..openrtc::client::BleConfig::default()
1904 }),
1905 ..openrtc::client::TransportConfig::default()
1906 }
1907 }
1908
1909 fn native_transport_install_context() -> NativeTransportInstallContext {
1910 NativeTransportInstallContext {
1911 data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
1912 }
1913 }
1914
1915 fn managed_session_test_result(user_id: &str) -> StartRtcManagedSessionResult {
1916 StartRtcManagedSessionResult {
1917 local_node_id: format!("node-{user_id}"),
1918 ticket_scope: Some("user-device".to_string()),
1919 ticket: None,
1920 presence_started: true,
1921 auto_connect_started: true,
1922 local_device: openrtc::native_device::NativeDeviceIdentity {
1923 device_id: format!("device-{user_id}"),
1924 device_name: "Test Device".to_string(),
1925 created_at_ms: 1,
1926 updated_at_ms: 1,
1927 name_source: None,
1928 system_info: None,
1929 },
1930 }
1931 }
1932
1933 async fn simulate_managed_session_start(
1934 state: Arc<OpenRtcTauriState>,
1935 api_key: String,
1936 user_id: String,
1937 pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
1938 starts: Arc<AtomicUsize>,
1939 stopped_owners: Arc<Mutex<Vec<usize>>>,
1940 ) -> ManagedSessionDisposition {
1941 let _start_guard = state.managed_session_start_guard.lock().await;
1942 let (client, _) = state.configure_identity(Some(api_key), None, None);
1943 let key = format!("{}:{user_id}:device-{user_id}", client.app_tag());
1944 let (disposition, active_epoch) = {
1945 let active = state.managed_session.lock().await;
1946 let disposition = managed_session_disposition(
1947 active.as_ref().map(|record| record.key.as_str()),
1948 active
1949 .as_ref()
1950 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1951 active
1952 .as_ref()
1953 .is_some_and(|record| record.result.presence_started),
1954 active
1955 .as_ref()
1956 .is_some_and(|record| record.result.auto_connect_started),
1957 &key,
1958 true,
1959 true,
1960 );
1961 (
1962 disposition,
1963 active.as_ref().map(|record| record.owner_epoch),
1964 )
1965 };
1966
1967 if disposition == ManagedSessionDisposition::Reuse {
1968 return disposition;
1969 }
1970
1971 if disposition == ManagedSessionDisposition::Replace {
1972 let previous = state
1973 .managed_session
1974 .lock()
1975 .await
1976 .take()
1977 .expect("replacement must have a previous owner");
1978 stopped_owners
1979 .lock()
1980 .await
1981 .push(Arc::as_ptr(&previous.owner_client) as usize);
1982 }
1983
1984 starts.fetch_add(1, Ordering::SeqCst);
1985 if let Some((entered, release)) = pause {
1986 entered.send(()).expect("start observer must be waiting");
1987 release.await.expect("start release must be sent");
1988 }
1989
1990 let owner_epoch = match disposition {
1991 ManagedSessionDisposition::Refresh => {
1992 active_epoch.expect("refresh must preserve the active owner epoch")
1993 }
1994 ManagedSessionDisposition::Start | ManagedSessionDisposition::Replace => {
1995 state.allocate_managed_session_owner_epoch()
1996 }
1997 ManagedSessionDisposition::Reuse => unreachable!("reuse returned before startup"),
1998 };
1999 let result = managed_session_test_result(&user_id);
2000 *state.managed_session.lock().await = Some(ManagedSessionRecord {
2001 key,
2002 owner_client: client,
2003 owner_epoch,
2004 result,
2005 });
2006 disposition
2007 }
2008
2009 #[tokio::test]
2010 async fn concurrent_same_key_start_reuses_one_owner_epoch() {
2011 let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
2012 let starts = Arc::new(AtomicUsize::new(0));
2013 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
2014 let api_key = "pk_test_same_managed_session_owner".to_string();
2015 let (first_entered_tx, first_entered_rx) = oneshot::channel();
2016 let (first_release_tx, first_release_rx) = oneshot::channel();
2017
2018 let first = tokio::spawn(simulate_managed_session_start(
2019 state.clone(),
2020 api_key.clone(),
2021 "same-user".to_string(),
2022 Some((first_entered_tx, first_release_rx)),
2023 starts.clone(),
2024 stopped_owners.clone(),
2025 ));
2026 first_entered_rx.await.expect("first start must pause");
2027
2028 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
2029 let second_state = state.clone();
2030 let second_starts = starts.clone();
2031 let second_stopped_owners = stopped_owners.clone();
2032 let second = tokio::spawn(async move {
2033 second_attempting_tx.send(()).unwrap();
2034 simulate_managed_session_start(
2035 second_state,
2036 api_key,
2037 "same-user".to_string(),
2038 None,
2039 second_starts,
2040 second_stopped_owners,
2041 )
2042 .await
2043 });
2044 second_attempting_rx.await.unwrap();
2045 tokio::task::yield_now().await;
2046 assert_eq!(starts.load(Ordering::SeqCst), 1);
2047
2048 first_release_tx.send(()).unwrap();
2049 assert_eq!(
2050 first.await.expect("first start task"),
2051 ManagedSessionDisposition::Start
2052 );
2053 assert_eq!(
2054 second.await.expect("second start task"),
2055 ManagedSessionDisposition::Reuse
2056 );
2057 assert_eq!(starts.load(Ordering::SeqCst), 1);
2058
2059 let active = state
2060 .managed_session
2061 .lock()
2062 .await
2063 .clone()
2064 .expect("active owner");
2065 assert_eq!(active.owner_epoch, 1);
2066 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
2067 assert!(stopped_owners.lock().await.is_empty());
2068 }
2069
2070 #[tokio::test]
2071 async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
2072 let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
2073 let starts = Arc::new(AtomicUsize::new(0));
2074 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
2075 let (first_entered_tx, first_entered_rx) = oneshot::channel();
2076 let (first_release_tx, first_release_rx) = oneshot::channel();
2077
2078 let first = tokio::spawn(simulate_managed_session_start(
2079 state.clone(),
2080 "pk_test_owner_a".to_string(),
2081 "first-user".to_string(),
2082 Some((first_entered_tx, first_release_rx)),
2083 starts.clone(),
2084 stopped_owners.clone(),
2085 ));
2086 first_entered_rx.await.expect("first start must pause");
2087 let first_owner = state.client();
2088
2089 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
2090 let second_state = state.clone();
2091 let second_starts = starts.clone();
2092 let second_stopped_owners = stopped_owners.clone();
2093 let second = tokio::spawn(async move {
2094 second_attempting_tx.send(()).unwrap();
2095 simulate_managed_session_start(
2096 second_state,
2097 "pk_test_owner_b".to_string(),
2098 "second-user".to_string(),
2099 None,
2100 second_starts,
2101 second_stopped_owners,
2102 )
2103 .await
2104 });
2105 second_attempting_rx.await.unwrap();
2106 tokio::task::yield_now().await;
2107 assert!(Arc::ptr_eq(&first_owner, &state.client()));
2108
2109 first_release_tx.send(()).unwrap();
2110 assert_eq!(
2111 first.await.expect("first start task"),
2112 ManagedSessionDisposition::Start
2113 );
2114 assert_eq!(
2115 second.await.expect("replacement start task"),
2116 ManagedSessionDisposition::Replace
2117 );
2118 assert_eq!(starts.load(Ordering::SeqCst), 2);
2119
2120 let active = state
2121 .managed_session
2122 .lock()
2123 .await
2124 .clone()
2125 .expect("active owner");
2126 assert_eq!(active.owner_epoch, 2);
2127 assert!(active.key.contains("second-user"));
2128 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
2129 assert_eq!(
2130 stopped_owners.lock().await.as_slice(),
2131 [Arc::as_ptr(&first_owner) as usize]
2132 );
2133 assert!(!Arc::ptr_eq(&active.owner_client, &first_owner));
2134 }
2135
2136 #[test]
2137 fn derives_app_tag_from_api_key_by_default() {
2138 let config = OpenRtcTauriConfig {
2139 api_key: Some("pk_test_1234567890abcdef".to_string()),
2140 ..OpenRtcTauriConfig::default()
2141 };
2142
2143 assert_eq!(config.app_tag(), "app_1234567890abcdef");
2144 }
2145
2146 #[tokio::test]
2147 async fn requested_native_transport_without_installer_keeps_base_route_available() {
2148 let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
2149 let client = state.client();
2150 ensure_requested_native_transports(
2151 &state,
2152 &client,
2153 Some(&requested_ble_config()),
2154 &native_transport_install_context(),
2155 )
2156 .await
2157 .expect("an unavailable optional transport must not fail base Iroh startup");
2158
2159 assert!(state.installed_native_transports.lock().await.is_empty());
2160 }
2161
2162 #[tokio::test]
2163 async fn native_transport_installer_is_idempotent_for_one_client() {
2164 let installs = Arc::new(AtomicUsize::new(0));
2165 let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
2166 .with_native_transport_installer(Arc::new(FakeNativeTransportInstaller {
2167 installs: installs.clone(),
2168 }));
2169 let client = state.client();
2170 let config = requested_ble_config();
2171
2172 ensure_requested_native_transports(
2173 &state,
2174 &client,
2175 Some(&config),
2176 &native_transport_install_context(),
2177 )
2178 .await
2179 .expect("first install");
2180 ensure_requested_native_transports(
2181 &state,
2182 &client,
2183 Some(&config),
2184 &native_transport_install_context(),
2185 )
2186 .await
2187 .expect("idempotent install");
2188
2189 assert_eq!(installs.load(Ordering::SeqCst), 1);
2190 }
2191
2192 #[tokio::test]
2193 async fn unavailable_native_transport_keeps_base_route_available() {
2194 let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
2195 .with_native_transport_installer(Arc::new(UnavailableNativeTransportInstaller));
2196 let client = state.client();
2197
2198 ensure_requested_native_transports(
2199 &state,
2200 &client,
2201 Some(&requested_ble_config()),
2202 &native_transport_install_context(),
2203 )
2204 .await
2205 .expect("optional transport installation failure must not fail base Iroh startup");
2206
2207 assert!(state.installed_native_transports.lock().await.is_empty());
2208 }
2209
2210 #[test]
2211 fn explicit_app_tag_wins_over_api_key() {
2212 let config = OpenRtcTauriConfig {
2213 api_key: Some("pk_test_1234567890abcdef".to_string()),
2214 app_tag: Some("space::manual".to_string()),
2215 ..OpenRtcTauriConfig::default()
2216 };
2217
2218 assert_eq!(config.app_tag(), "space::manual");
2219 }
2220
2221 #[test]
2222 fn derives_space_app_tag_from_space_key() {
2223 let config = OpenRtcTauriConfig {
2224 api_key: Some("pk_test_abc".to_string()),
2225 space_key: Some("room".to_string()),
2226 ..OpenRtcTauriConfig::default()
2227 };
2228
2229 assert!(config.app_tag().starts_with("space::"));
2230 assert_eq!(config.app_tag().len(), "space::".len() + 64);
2231 }
2232
2233 #[test]
2234 fn from_env_accepts_vite_project_and_app_tag() {
2235 let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
2236 let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
2237 std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
2238 std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
2239
2240 let config = OpenRtcTauriConfig::from_env();
2241
2242 assert_eq!(config.project_id, "pluto-rtc-prod");
2243 assert_eq!(config.app_tag(), "app_from_vite_env");
2244
2245 match previous_project {
2246 Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
2247 None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
2248 }
2249 match previous_app_tag {
2250 Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
2251 None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
2252 }
2253 }
2254
2255 #[test]
2256 fn native_state_reconfigures_anonymous_client_to_public_space() {
2257 let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
2258 assert_eq!(state.client().app_tag(), "app_anonymous");
2259
2260 let (client, reconfigured) = state.configure_identity(
2261 Some("pk_test_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string()),
2262 Some("assistant-control-plane".to_string()),
2263 None,
2264 );
2265
2266 assert!(reconfigured);
2267 assert!(client.app_tag().starts_with("space::"));
2268 assert_eq!(client.app_tag(), state.client().app_tag());
2269 assert_ne!(client.app_tag(), "app_anonymous");
2270 }
2271
2272 #[test]
2273 fn managed_session_uses_persisted_native_device_id() {
2274 let identity = openrtc::native_device::NativeDeviceIdentity {
2275 device_id: "persisted-native-device".to_string(),
2276 device_name: "Mac".to_string(),
2277 created_at_ms: 1,
2278 updated_at_ms: 1,
2279 name_source: None,
2280 system_info: None,
2281 };
2282
2283 assert_eq!(
2284 managed_session_device_id(Some(" desktop-e2e-native "), &identity),
2285 "persisted-native-device"
2286 );
2287 assert_eq!(
2288 managed_session_device_id(None, &identity),
2289 "persisted-native-device"
2290 );
2291 }
2292
2293 #[test]
2294 fn repeated_managed_session_triggers_converge_on_one_owner() {
2295 let key = "app:user:persisted-native-device";
2296 let cases = [
2297 ("react-remount", true, true, true, true),
2298 ("hmr", true, true, true, true),
2299 ("auth-refresh", true, true, true, true),
2300 ("resume", true, true, true, true),
2301 ("alias-change", true, true, true, true),
2302 ];
2303 for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
2304 assert_eq!(
2305 managed_session_disposition(
2306 Some(key),
2307 true,
2308 active_presence,
2309 active_auto,
2310 key,
2311 requested_presence,
2312 requested_auto,
2313 ),
2314 ManagedSessionDisposition::Reuse,
2315 "{label} must reuse the authoritative tuple"
2316 );
2317 }
2318
2319 assert_eq!(
2320 managed_session_disposition(Some(key), true, false, true, key, true, true),
2321 ManagedSessionDisposition::Refresh,
2322 "failed presence startup must retry idempotently"
2323 );
2324 assert_eq!(
2325 managed_session_disposition(
2326 Some(key),
2327 true,
2328 true,
2329 true,
2330 "app:other-user:persisted-native-device",
2331 true,
2332 true,
2333 ),
2334 ManagedSessionDisposition::Replace,
2335 "user switching must replace the previous lifecycle owner"
2336 );
2337 assert_eq!(
2338 managed_session_disposition(Some(key), false, true, true, key, true, true),
2339 ManagedSessionDisposition::Replace,
2340 "the same tuple on a different client epoch must replace the previous owner"
2341 );
2342 }
2343
2344 #[test]
2345 fn managed_session_metadata_uses_authoritative_device_id() {
2346 let raw = serde_json::json!({
2347 "deviceId": "persisted-native-device",
2348 "assistantDevice": {
2349 "deviceId": "desktop-e2e-native"
2350 }
2351 })
2352 .to_string();
2353
2354 let metadata =
2355 metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
2356 let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
2357
2358 assert_eq!(parsed["deviceId"], "desktop-e2e-native");
2359 assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
2360 }
2361}