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