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