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