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