Skip to main content

openrtc_tauri_plugin/
lib.rs

1use std::collections::HashMap;
2use std::path::PathBuf;
3use std::sync::{Arc, RwLock};
4
5use futures::StreamExt;
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 RTC_DEVICE_EVENTS: &str = "rtc-device-events";
14const RTC_SESSION_EVENTS: &str = "rtc-session-events";
15const CONNECTION_STATE_CHANGED_EVENT: &str = "connection-state-changed";
16const PEER_BI_STREAM_INCOMING_EVENT: &str = "openrtc://peer-bi-stream/incoming";
17const PEER_BI_STREAM_CHUNK_EVENT: &str = "openrtc://peer-bi-stream/chunk";
18const PEER_BI_STREAM_CLOSED_EVENT: &str = "openrtc://peer-bi-stream/closed";
19
20#[derive(Debug, Clone, Serialize, Deserialize)]
21#[serde(rename_all = "camelCase")]
22pub struct OpenRtcTauriConfig {
23    #[serde(default = "default_project_id")]
24    pub project_id: String,
25    #[serde(default)]
26    pub api_key: Option<String>,
27    #[serde(default)]
28    pub app_tag: Option<String>,
29    #[serde(default)]
30    pub space_key: Option<String>,
31    #[serde(default)]
32    pub data_dir: Option<PathBuf>,
33    #[serde(default)]
34    pub transport_config: Option<openrtc::client::TransportConfig>,
35}
36
37impl Default for OpenRtcTauriConfig {
38    fn default() -> Self {
39        Self {
40            project_id: default_project_id(),
41            api_key: None,
42            app_tag: None,
43            space_key: None,
44            data_dir: None,
45            transport_config: None,
46        }
47    }
48}
49
50impl OpenRtcTauriConfig {
51    pub fn from_env() -> Self {
52        let api_key = first_env(&[
53            "VITE_OPENRTC_KEY",
54            "VITE_OPENRTC_API_KEY",
55            "VITE_PLUTO_OPENRTC_API_KEY",
56            "OPENRTC_API_KEY",
57        ]);
58        let space_key = first_env(&[
59            "VITE_OPENRTC_SPACE_KEY",
60            "VITE_OPENRTC_SPACE",
61            "OPENRTC_SPACE_KEY",
62        ]);
63        let app_tag = first_env(&["OPENRTC_APP_TAG"]);
64        let project_id = first_env(&["OPENRTC_PROJECT_ID"]).unwrap_or_else(default_project_id);
65
66        Self {
67            project_id,
68            api_key,
69            app_tag,
70            space_key,
71            data_dir: None,
72            transport_config: None,
73        }
74    }
75
76    pub fn app_tag(&self) -> String {
77        if let Some(app_tag) = self
78            .app_tag
79            .as_deref()
80            .map(str::trim)
81            .filter(|value| !value.is_empty())
82        {
83            return app_tag.to_string();
84        }
85
86        if let (Some(api_key), Some(space_key)) = (
87            self.api_key
88                .as_deref()
89                .map(str::trim)
90                .filter(|value| !value.is_empty()),
91            self.space_key
92                .as_deref()
93                .map(str::trim)
94                .filter(|value| !value.is_empty()),
95        ) {
96            return openrtc::space_app_tag_from_keys(api_key, space_key);
97        }
98
99        self.api_key
100            .as_deref()
101            .map(openrtc::app_tag_from_api_key)
102            .unwrap_or_else(|| "app_anonymous".to_string())
103    }
104}
105
106fn default_project_id() -> String {
107    openrtc::LIVE_PROJECT_ID.to_string()
108}
109
110fn first_env(names: &[&str]) -> Option<String> {
111    names.iter().find_map(|name| {
112        std::env::var(name)
113            .ok()
114            .map(|value| value.trim().to_string())
115            .filter(|value| !value.is_empty())
116    })
117}
118
119#[derive(Default)]
120struct TokenRelayState {
121    auth_token: RwLock<Option<String>>,
122    refresh_token: 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 || relay.auth_token.read().ok().and_then(|guard| guard.clone()))
129    }
130
131    fn set(&self, auth_token: Option<String>, refresh_token: Option<String>) {
132        if let Ok(mut guard) = self.auth_token.write() {
133            *guard = normalize_token(auth_token);
134        }
135        if let Ok(mut guard) = self.refresh_token.write() {
136            *guard = normalize_token(refresh_token);
137        }
138    }
139}
140
141fn normalize_token(value: Option<String>) -> Option<String> {
142    value
143        .map(|value| value.trim().to_string())
144        .filter(|value| !value.is_empty())
145}
146
147pub struct OpenRtcTauriState {
148    client: RwLock<Arc<openrtc::client::Client>>,
149    config: RwLock<OpenRtcTauriConfig>,
150    token_relay: Arc<TokenRelayState>,
151    data_dir: Option<PathBuf>,
152    connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
153    loop_state: Mutex<Option<PresenceLoopParams>>,
154    subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
155    peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
156}
157
158impl OpenRtcTauriState {
159    pub fn new(config: OpenRtcTauriConfig) -> Self {
160        openrtc::ensure_default_rustls_provider();
161
162        let token_relay = Arc::new(TokenRelayState::default());
163        let client = build_client(&config, &token_relay);
164        let data_dir = config.data_dir.clone();
165
166        Self {
167            client: RwLock::new(client),
168            config: RwLock::new(config),
169            token_relay,
170            data_dir,
171            connection_state_forwarder: std::sync::Mutex::new(None),
172            loop_state: Mutex::new(None),
173            subscriptions: Mutex::new(HashMap::new()),
174            peer_bi_streams: Mutex::new(HashMap::new()),
175        }
176    }
177
178    pub fn client(&self) -> Arc<openrtc::client::Client> {
179        self.client
180            .read()
181            .map(|guard| guard.clone())
182            .unwrap_or_else(|_| build_client(&OpenRtcTauriConfig::default(), &self.token_relay))
183    }
184
185    pub fn set_auth_token(&self, auth_token: Option<String>, refresh_token: Option<String>) {
186        self.token_relay.set(auth_token, refresh_token);
187    }
188
189    fn configure_identity(
190        &self,
191        api_key: Option<String>,
192        space_key: Option<String>,
193        app_tag: Option<String>,
194    ) -> (Arc<openrtc::client::Client>, bool) {
195        let mut next_config = self
196            .config
197            .read()
198            .map(|guard| guard.clone())
199            .unwrap_or_default();
200        apply_identity_override(&mut next_config.api_key, api_key);
201        apply_identity_override(&mut next_config.space_key, space_key);
202        apply_identity_override(&mut next_config.app_tag, app_tag);
203
204        let next_app_tag = next_config.app_tag();
205        let current = self.client();
206        if current.app_tag() == next_app_tag {
207            if let Ok(mut guard) = self.config.write() {
208                *guard = next_config;
209            }
210            return (current, false);
211        }
212
213        let next_client = build_client(&next_config, &self.token_relay);
214        if let Ok(mut guard) = self.client.write() {
215            *guard = next_client.clone();
216        }
217        if let Ok(mut guard) = self.config.write() {
218            *guard = next_config;
219        }
220        // Reconfiguration changes the app identity boundary. Stop scoped
221        // background work on the old client after replacing it in state.
222        current.stop_auth_scoped_activity();
223        (next_client, true)
224    }
225
226    fn replace_connection_state_forwarder<R: Runtime>(
227        &self,
228        app: tauri::AppHandle<R>,
229        client: Arc<openrtc::client::Client>,
230    ) {
231        let next = tauri::async_runtime::spawn(async move {
232            forward_connection_state_events(app, client).await;
233        });
234        if let Ok(mut guard) = self.connection_state_forwarder.lock() {
235            if let Some(previous) = guard.replace(next) {
236                previous.abort();
237            }
238        }
239    }
240}
241
242fn build_client(
243    config: &OpenRtcTauriConfig,
244    token_relay: &Arc<TokenRelayState>,
245) -> Arc<openrtc::client::Client> {
246    let mut builder = openrtc::client::Client::builder_with_app_tag(
247        config.project_id.clone(),
248        config.app_tag(),
249        token_relay.token_provider(),
250    );
251    if let Some(transport_config) = config.transport_config.clone() {
252        builder = builder.transport_config(transport_config);
253    }
254    Arc::new(builder.build())
255}
256
257fn apply_identity_override(target: &mut Option<String>, value: Option<String>) {
258    if let Some(value) = value
259        .as_deref()
260        .map(str::trim)
261        .filter(|value| !value.is_empty())
262    {
263        *target = Some(value.to_string());
264    }
265}
266
267fn requested_local_device_id(value: Option<&str>) -> Option<String> {
268    value
269        .map(str::trim)
270        .filter(|value| !value.is_empty())
271        .map(ToOwned::to_owned)
272}
273
274fn managed_session_device_id(
275    requested: Option<&str>,
276    native_identity: &openrtc::native_device::NativeDeviceIdentity,
277) -> String {
278    requested_local_device_id(requested).unwrap_or_else(|| native_identity.device_id.clone())
279}
280
281fn metadata_with_authoritative_device_id(
282    metadata: Option<String>,
283    device_id: &str,
284) -> Option<String> {
285    let device_id = device_id.trim();
286    if device_id.is_empty() {
287        return metadata;
288    }
289
290    let Some(raw_metadata) = metadata else {
291        return Some(serde_json::json!({ "deviceId": device_id }).to_string());
292    };
293
294    match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
295        Ok(serde_json::Value::Object(mut map)) => {
296            map.insert(
297                "deviceId".to_string(),
298                serde_json::Value::String(device_id.to_string()),
299            );
300            Some(serde_json::Value::Object(map).to_string())
301        }
302        _ => Some(
303            serde_json::json!({
304                "deviceId": device_id,
305                "metadata": raw_metadata,
306            })
307            .to_string(),
308        ),
309    }
310}
311
312#[derive(Debug, Clone)]
313struct PresenceLoopParams {
314    user_id: String,
315    device_name: String,
316    ticket: String,
317    metadata: Option<String>,
318}
319
320struct PeerBiStreamHandle {
321    send: Arc<Mutex<Option<PeerSendStream>>>,
322    recv: Option<PeerRecvStream>,
323    read_task: Option<tokio::task::JoinHandle<()>>,
324}
325
326#[derive(Serialize)]
327#[serde(rename_all = "camelCase")]
328pub struct OpenPeerBiStreamResult {
329    stream_id: String,
330    connection_id: Option<String>,
331    remote_node_id: String,
332}
333
334#[derive(Serialize, Clone)]
335#[serde(rename_all = "camelCase")]
336struct IncomingPeerBiStreamEvent {
337    request_id: String,
338    stream_id: String,
339    remote_node_id: String,
340}
341
342#[derive(Serialize, Clone)]
343#[serde(rename_all = "camelCase")]
344struct PeerBiStreamChunkEvent {
345    stream_id: String,
346    bytes: Vec<u8>,
347}
348
349#[derive(Serialize, Clone)]
350#[serde(rename_all = "camelCase")]
351struct PeerBiStreamClosedEvent {
352    stream_id: String,
353    error: Option<String>,
354}
355
356#[derive(Debug, Clone, Serialize)]
357#[serde(rename_all = "camelCase")]
358pub struct SearchRtcDevicesResult {
359    devices: Vec<openrtc::signaling::Device>,
360}
361
362#[derive(Debug, Clone, Serialize)]
363#[serde(rename_all = "camelCase")]
364pub struct SearchRtcDevicesWithStatusResult {
365    devices: Vec<openrtc::client::DeviceStatusSnapshot>,
366}
367
368#[derive(Debug, Clone, Serialize)]
369#[serde(rename_all = "camelCase")]
370pub struct StartRtcManagedSessionResult {
371    local_node_id: String,
372    ticket_scope: Option<String>,
373    ticket: Option<String>,
374    presence_started: bool,
375    auto_connect_started: bool,
376    local_device: openrtc::native_device::NativeDeviceIdentity,
377}
378
379pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
380    init_with_state(OpenRtcTauriState::new(config))
381}
382
383pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
384    tauri::plugin::Builder::new(PLUGIN_NAME)
385        .setup(move |app, _api| {
386            let client = state.client();
387            state.replace_connection_state_forwarder(app.clone(), client);
388            app.manage(state);
389            Ok(())
390        })
391        .invoke_handler(tauri::generate_handler![
392            desktop_set_pluto_auth_token,
393            openrtc_set_auth_token,
394            rtc_native_status,
395            get_rtc_local_device_info,
396            update_rtc_local_device_name,
397            get_iroh_node_id,
398            start_iroh_node,
399            get_iroh_endpoint_ticket,
400            register_session_token,
401            get_endpoint_ticket_with_token,
402            validate_session_token,
403            revoke_session_tokens_by_scope,
404            search_rtc_devices,
405            search_rtc_devices_with_status,
406            update_rtc_device,
407            delete_rtc_device,
408            set_rtc_offline,
409            start_rtc_presence_loop,
410            stop_rtc_presence_loop,
411            start_rtc_managed_session,
412            start_rtc_auto_connect,
413            stop_rtc_auto_connect,
414            force_rtc_reconnect_snapshot,
415            connect_to_device,
416            disconnect_device,
417            set_auto_connect_excluded,
418            resolve_rtc_peer_connection_records,
419            resolve_rtc_peer_identity,
420            get_rtc_peer_session,
421            list_rtc_peer_sessions,
422            list_rtc_managed_connections,
423            wait_for_rtc_settled_peer,
424            list_rtc_connection_states,
425            get_rtc_connection_state,
426            start_rtc_device_subscription,
427            start_rtc_session_subscription,
428            stop_rtc_subscription,
429            start_incoming_peer_bi_streams,
430            open_peer_bi_stream,
431            open_peer_bi_transport_only_stream,
432            open_peer_native_bi_stream,
433            write_peer_bi_stream,
434            start_peer_bi_stream_read,
435            close_peer_bi_stream,
436            get_app_limits,
437        ])
438        .build()
439}
440
441pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
442    init(OpenRtcTauriConfig::from_env())
443}
444
445async fn forward_connection_state_events<R: Runtime>(
446    app: tauri::AppHandle<R>,
447    client: Arc<openrtc::client::Client>,
448) {
449    let mut rx = client.subscribe_native_connection_state_updates();
450    while let Ok(snapshot) = rx.recv().await {
451        let _ = app.emit(CONNECTION_STATE_CHANGED_EVENT, snapshot);
452    }
453}
454
455fn app_data_dir<R: Runtime>(
456    app: &tauri::AppHandle<R>,
457    state: &OpenRtcTauriState,
458) -> Result<PathBuf, String> {
459    state
460        .data_dir
461        .clone()
462        .or_else(|| app.path().app_data_dir().ok())
463        .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
464}
465
466async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
467    if let Some(node_id) = client.current_node_id().await {
468        return Ok(node_id);
469    }
470    client
471        .init_iroh(None, Vec::new())
472        .await
473        .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
474}
475
476async fn register_peer_bi_stream<R: Runtime>(
477    app: tauri::AppHandle<R>,
478    state: &OpenRtcTauriState,
479    connection_id: Option<String>,
480    remote_node_id: String,
481    send: PeerSendStream,
482    recv: PeerRecvStream,
483    start_reading: bool,
484) -> OpenPeerBiStreamResult {
485    let stream_id = uuid::Uuid::new_v4().to_string();
486    let send = Arc::new(Mutex::new(Some(send)));
487    let (recv, read_task) = if start_reading {
488        (
489            None,
490            Some(spawn_peer_bi_stream_reader(
491                app.clone(),
492                stream_id.clone(),
493                recv,
494            )),
495        )
496    } else {
497        // Incoming native streams are announced to JS before forwarding bytes so
498        // the bridge can install per-stream listeners without losing the first
499        // native-main label/session-token frame.
500        (Some(recv), None)
501    };
502
503    state.peer_bi_streams.lock().await.insert(
504        stream_id.clone(),
505        PeerBiStreamHandle {
506            send,
507            recv,
508            read_task,
509        },
510    );
511
512    OpenPeerBiStreamResult {
513        stream_id,
514        connection_id,
515        remote_node_id,
516    }
517}
518
519fn spawn_peer_bi_stream_reader<R: Runtime>(
520    app: tauri::AppHandle<R>,
521    stream_id: String,
522    mut recv: PeerRecvStream,
523) -> tokio::task::JoinHandle<()> {
524    tokio::spawn(async move {
525        let mut chunk = vec![0_u8; 64 * 1024];
526        let mut close_error: Option<String> = None;
527        loop {
528            match recv.read(&mut chunk).await {
529                Ok(0) => break,
530                Ok(n) => {
531                    let _ = app.emit(
532                        PEER_BI_STREAM_CHUNK_EVENT,
533                        PeerBiStreamChunkEvent {
534                            stream_id: stream_id.clone(),
535                            bytes: chunk[..n].to_vec(),
536                        },
537                    );
538                }
539                Err(error) => {
540                    close_error = Some(error.to_string());
541                    break;
542                }
543            }
544        }
545        let _ = app.emit(
546            PEER_BI_STREAM_CLOSED_EVENT,
547            PeerBiStreamClosedEvent {
548                stream_id,
549                error: close_error,
550            },
551        );
552    })
553}
554
555async fn open_peer_bi_with<R: Runtime, F, Fut>(
556    app: tauri::AppHandle<R>,
557    state: tauri::State<'_, OpenRtcTauriState>,
558    peer_id: String,
559    timeout_ms: Option<u64>,
560    open: F,
561) -> Result<OpenPeerBiStreamResult, String>
562where
563    F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
564    Fut: std::future::Future<
565        Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
566    >,
567{
568    let peer_id = peer_id.trim().to_string();
569    if peer_id.is_empty() {
570        return Err("peerId is required".to_string());
571    }
572
573    let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
574        .await
575        .map_err(|error| format!("open peer bi stream failed: {error}"))?;
576    Ok(register_peer_bi_stream(app, &state, connection_id, remote_node_id, send, recv, true).await)
577}
578
579#[tauri::command]
580async fn desktop_set_pluto_auth_token(
581    state: tauri::State<'_, OpenRtcTauriState>,
582    auth_token: Option<String>,
583    refresh_token: Option<String>,
584) -> Result<(), String> {
585    state.set_auth_token(auth_token, refresh_token);
586    if let Some(params) = state.loop_state.lock().await.clone() {
587        let client = state.client();
588        tokio::spawn(async move {
589            let _ = client
590                .update_presence(
591                    &params.user_id,
592                    &params.device_name,
593                    &params.ticket,
594                    params.metadata.as_deref(),
595                )
596                .await;
597        });
598    }
599    Ok(())
600}
601
602#[tauri::command]
603async fn openrtc_set_auth_token(
604    state: tauri::State<'_, OpenRtcTauriState>,
605    auth_token: Option<String>,
606    refresh_token: Option<String>,
607) -> Result<(), String> {
608    desktop_set_pluto_auth_token(state, auth_token, refresh_token).await
609}
610
611#[tauri::command]
612async fn rtc_native_status(
613    state: tauri::State<'_, OpenRtcTauriState>,
614) -> Result<openrtc::client::RuntimeStatus, String> {
615    Ok(state.client().runtime_status().await)
616}
617
618#[tauri::command]
619async fn get_rtc_local_device_info<R: Runtime>(
620    app: tauri::AppHandle<R>,
621    state: tauri::State<'_, OpenRtcTauriState>,
622) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
623    let data_dir = app_data_dir(&app, &state)?;
624    state
625        .client()
626        .init_native_device_identity(data_dir, None)
627        .await
628        .map_err(|error| error.to_string())
629}
630
631#[tauri::command]
632async fn update_rtc_local_device_name<R: Runtime>(
633    app: tauri::AppHandle<R>,
634    state: tauri::State<'_, OpenRtcTauriState>,
635    device_name: String,
636) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
637    let data_dir = app_data_dir(&app, &state)?;
638    let _ = state
639        .client()
640        .init_native_device_identity(data_dir, None)
641        .await
642        .map_err(|error| error.to_string())?;
643    state
644        .client()
645        .update_native_device_name(&device_name)
646        .await
647        .map_err(|error| error.to_string())
648}
649
650#[tauri::command]
651async fn get_iroh_node_id(
652    state: tauri::State<'_, OpenRtcTauriState>,
653) -> Result<Option<String>, String> {
654    Ok(state.client().current_node_id().await)
655}
656
657#[tauri::command]
658async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
659    ensure_iroh_node(state.client()).await
660}
661
662#[tauri::command]
663async fn get_iroh_endpoint_ticket(
664    state: tauri::State<'_, OpenRtcTauriState>,
665) -> Result<String, String> {
666    ensure_iroh_node(state.client()).await?;
667    state
668        .client()
669        .endpoint_ticket()
670        .await
671        .map_err(|error| error.to_string())
672}
673
674#[tauri::command]
675async fn register_session_token(
676    state: tauri::State<'_, OpenRtcTauriState>,
677    token: String,
678    scope: String,
679    max_connections: u32,
680    expires_at_ms: Option<u64>,
681) -> Result<(), String> {
682    if let Some(expires_at_ms) = expires_at_ms {
683        state.client().register_session_token_with_expiry_ms(
684            token,
685            scope,
686            max_connections,
687            expires_at_ms,
688        );
689    } else {
690        state
691            .client()
692            .register_session_token(token, scope, max_connections);
693    }
694    Ok(())
695}
696
697#[tauri::command]
698async fn get_endpoint_ticket_with_token(
699    state: tauri::State<'_, OpenRtcTauriState>,
700    scope: String,
701    max_connections: u32,
702) -> Result<String, String> {
703    ensure_iroh_node(state.client()).await?;
704    state
705        .client()
706        .endpoint_ticket_with_token(&scope, max_connections)
707        .await
708        .map_err(|error| error.to_string())
709}
710
711#[tauri::command]
712async fn validate_session_token(
713    state: tauri::State<'_, OpenRtcTauriState>,
714    token: String,
715    connection_id: Option<String>,
716) -> Result<String, String> {
717    if let Some(connection_id) = connection_id
718        .as_deref()
719        .map(str::trim)
720        .filter(|value| !value.is_empty())
721    {
722        state
723            .client()
724            .validate_session_token_for_connection(&token, connection_id)
725            .await
726            .map_err(|error| error.to_string())
727    } else {
728        state.client().validate_session_token(&token)
729    }
730}
731
732#[tauri::command]
733async fn revoke_session_tokens_by_scope(
734    state: tauri::State<'_, OpenRtcTauriState>,
735    scope: String,
736) -> Result<Vec<String>, String> {
737    Ok(state.client().revoke_tokens_by_scope(&scope).await)
738}
739
740#[tauri::command]
741async fn search_rtc_devices(
742    state: tauri::State<'_, OpenRtcTauriState>,
743    user_id: String,
744) -> Result<SearchRtcDevicesResult, String> {
745    let devices = state
746        .client()
747        .search_devices(&user_id)
748        .await
749        .map_err(|error| error.to_string())?;
750    Ok(SearchRtcDevicesResult { devices })
751}
752
753#[tauri::command]
754async fn search_rtc_devices_with_status(
755    state: tauri::State<'_, OpenRtcTauriState>,
756    user_id: String,
757) -> Result<SearchRtcDevicesWithStatusResult, String> {
758    let devices = state
759        .client()
760        .devices_with_status(&user_id)
761        .await
762        .map_err(|error| error.to_string())?;
763    Ok(SearchRtcDevicesWithStatusResult { devices })
764}
765
766#[tauri::command]
767async fn update_rtc_device(
768    state: tauri::State<'_, OpenRtcTauriState>,
769    user_id: String,
770    device_id: String,
771    device_name: Option<String>,
772    capabilities: Option<openrtc::signaling::DeviceCapabilities>,
773    metadata: Option<String>,
774) -> Result<(), String> {
775    state
776        .client()
777        .update_device(
778            &user_id,
779            &device_id,
780            device_name.as_deref(),
781            capabilities,
782            metadata.as_deref(),
783        )
784        .await
785        .map_err(|error| error.to_string())
786}
787
788#[tauri::command]
789async fn delete_rtc_device(
790    state: tauri::State<'_, OpenRtcTauriState>,
791    user_id: String,
792    device_id: String,
793) -> Result<(), String> {
794    state
795        .client()
796        .delete_device(&user_id, &device_id)
797        .await
798        .map_err(|error| error.to_string())
799}
800
801#[tauri::command]
802async fn set_rtc_offline(
803    state: tauri::State<'_, OpenRtcTauriState>,
804    user_id: String,
805) -> Result<(), String> {
806    state
807        .client()
808        .set_offline(&user_id)
809        .await
810        .map_err(|error| error.to_string())
811}
812
813#[tauri::command]
814async fn start_rtc_presence_loop(
815    state: tauri::State<'_, OpenRtcTauriState>,
816    user_id: String,
817    local_node_id: String,
818    device_name: String,
819    ticket: String,
820    metadata: Option<String>,
821    api_key: Option<String>,
822    space_key: Option<String>,
823) -> Result<(), String> {
824    let _ = local_node_id;
825    let (client, _) = state.configure_identity(api_key.clone(), space_key, None);
826    if let Some(api_key) = api_key
827        .as_deref()
828        .map(str::trim)
829        .filter(|value| !value.is_empty())
830    {
831        let api_key = api_key.to_string();
832        let client = client.clone();
833        tokio::spawn(async move {
834            let _ = client.fetch_and_cache_app_limits(&api_key).await;
835        });
836    }
837    client.start_signaling_loop(
838        user_id.clone(),
839        device_name.clone(),
840        ticket.clone(),
841        metadata.clone(),
842    );
843    *state.loop_state.lock().await = Some(PresenceLoopParams {
844        user_id,
845        device_name,
846        ticket,
847        metadata,
848    });
849    Ok(())
850}
851
852#[tauri::command]
853async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
854    state.client().stop_presence_loop();
855    *state.loop_state.lock().await = None;
856    stop_subscription_by_id(&state, "internal-presence-loop").await;
857    Ok(())
858}
859
860#[tauri::command]
861async fn start_rtc_managed_session<R: Runtime>(
862    app: tauri::AppHandle<R>,
863    state: tauri::State<'_, OpenRtcTauriState>,
864    user_id: String,
865    device_name: Option<String>,
866    local_device_id: Option<String>,
867    metadata: Option<String>,
868    transports: Option<openrtc::client::TransportConfig>,
869    auto_connect: Option<bool>,
870    presence: Option<bool>,
871    api_key: Option<String>,
872    space_key: Option<String>,
873) -> Result<StartRtcManagedSessionResult, String> {
874    let (client, reconfigured) = state.configure_identity(api_key.clone(), space_key.clone(), None);
875    if reconfigured {
876        state.replace_connection_state_forwarder(app.clone(), client.clone());
877        eprintln!(
878            "[openrtc-tauri] reconfigured native OpenRTC namespace app_tag={}",
879            client.app_tag()
880        );
881    }
882
883    if let Some(transport_config) = transports {
884        client
885            .update_transport_config(transport_config)
886            .await
887            .map_err(|error| error.to_string())?;
888    }
889
890    let local_node_id = ensure_iroh_node(client.clone()).await?;
891    let data_dir = app_data_dir(&app, &state)?;
892    let local_device = client
893        .init_native_device_identity(data_dir, device_name.as_deref())
894        .await
895        .map_err(|error| error.to_string())?;
896
897    let effective_local_device_id =
898        managed_session_device_id(local_device_id.as_deref(), &local_device);
899    if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
900        if requested != local_device.device_id {
901            eprintln!(
902                "[openrtc-tauri] requested localDeviceId={} using logical signaling identity; native runtime identity is {}",
903                requested, local_device.device_id
904            );
905        }
906    }
907
908    let should_start_presence = presence.unwrap_or(true);
909    let should_start_auto_connect = auto_connect.unwrap_or(true);
910    let ticket = if should_start_presence || should_start_auto_connect {
911        Some(
912            client
913                .endpoint_ticket_with_token("user-device", 0)
914                .await
915                .map_err(|error| error.to_string())?,
916        )
917    } else {
918        None
919    };
920    let managed_metadata =
921        metadata_with_authoritative_device_id(metadata.clone(), &effective_local_device_id);
922
923    if should_start_presence {
924        start_rtc_presence_loop(
925            state.clone(),
926            user_id.clone(),
927            local_node_id.clone(),
928            local_device.device_name.clone(),
929            ticket
930                .clone()
931                .ok_or_else(|| "managed session expected a ticket".to_string())?,
932            managed_metadata.clone(),
933            api_key,
934            space_key,
935        )
936        .await?;
937    }
938
939    if should_start_auto_connect {
940        client.start_auto_connect(user_id.clone(), effective_local_device_id.clone());
941    }
942
943    let mut logical_local_device = local_device.clone();
944    logical_local_device.device_id = effective_local_device_id;
945
946    Ok(StartRtcManagedSessionResult {
947        local_node_id,
948        ticket_scope: ticket.as_ref().map(|_| "user-device".to_string()),
949        ticket,
950        presence_started: should_start_presence,
951        auto_connect_started: should_start_auto_connect,
952        local_device: logical_local_device,
953    })
954}
955
956#[tauri::command]
957async fn start_rtc_auto_connect(
958    state: tauri::State<'_, OpenRtcTauriState>,
959    user_id: String,
960    local_device_id: String,
961) -> Result<(), String> {
962    state
963        .client()
964        .clone()
965        .start_auto_connect(user_id, local_device_id);
966    Ok(())
967}
968
969#[tauri::command]
970async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
971    state.client().stop_auto_connect();
972    Ok(())
973}
974
975#[tauri::command]
976async fn force_rtc_reconnect_snapshot(
977    state: tauri::State<'_, OpenRtcTauriState>,
978) -> Result<(), String> {
979    state.client().clone().force_reconnect_snapshot();
980    Ok(())
981}
982
983#[tauri::command]
984async fn connect_to_device(
985    state: tauri::State<'_, OpenRtcTauriState>,
986    device_id: Option<String>,
987    endpoint_ticket: String,
988) -> Result<openrtc::client::ManagedConnectResult, String> {
989    state
990        .client()
991        .connect_device(device_id.as_deref(), &endpoint_ticket)
992        .await
993        .map_err(|error| error.to_string())
994}
995
996#[tauri::command]
997async fn disconnect_device(
998    state: tauri::State<'_, OpenRtcTauriState>,
999    device_id: String,
1000    node_id_hint: Option<String>,
1001) -> Result<(), String> {
1002    state
1003        .client()
1004        .disconnect_device(&device_id, node_id_hint.as_deref())
1005        .await;
1006    Ok(())
1007}
1008
1009#[tauri::command]
1010async fn set_auto_connect_excluded(
1011    state: tauri::State<'_, OpenRtcTauriState>,
1012    device_id: String,
1013    excluded: bool,
1014) -> Result<(), String> {
1015    if excluded {
1016        state.client().exclude_peer_and_publish(&device_id).await;
1017    } else {
1018        state.client().unexclude_peer_and_publish(&device_id).await;
1019    }
1020    Ok(())
1021}
1022
1023#[tauri::command]
1024async fn resolve_rtc_peer_connection_records(
1025    state: tauri::State<'_, OpenRtcTauriState>,
1026    id: String,
1027) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1028    Ok(state.client().resolve_peer_connection_records(&id).await)
1029}
1030
1031#[tauri::command]
1032async fn resolve_rtc_peer_identity(
1033    state: tauri::State<'_, OpenRtcTauriState>,
1034    id: String,
1035) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
1036    Ok(state.client().peer_snapshot(&id).await)
1037}
1038
1039#[tauri::command]
1040async fn get_rtc_peer_session(
1041    state: tauri::State<'_, OpenRtcTauriState>,
1042    id: String,
1043) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1044    Ok(state.client().peer_session(&id).await)
1045}
1046
1047#[tauri::command]
1048async fn list_rtc_peer_sessions(
1049    state: tauri::State<'_, OpenRtcTauriState>,
1050) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
1051    Ok(state.client().peer_sessions().await)
1052}
1053
1054#[tauri::command]
1055async fn list_rtc_managed_connections(
1056    state: tauri::State<'_, OpenRtcTauriState>,
1057) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1058    Ok(state.client().list_managed_connections().await)
1059}
1060
1061#[tauri::command]
1062async fn wait_for_rtc_settled_peer(
1063    state: tauri::State<'_, OpenRtcTauriState>,
1064    id: String,
1065    timeout_ms: Option<u64>,
1066) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1067    Ok(state.client().wait_for_settled_peer(&id, timeout_ms).await)
1068}
1069
1070#[tauri::command]
1071async fn list_rtc_connection_states(
1072    state: tauri::State<'_, OpenRtcTauriState>,
1073) -> Result<Vec<openrtc::client::ConnectionStateSnapshot>, String> {
1074    Ok(state.client().connection_states().await)
1075}
1076
1077#[tauri::command]
1078async fn get_rtc_connection_state(
1079    state: tauri::State<'_, OpenRtcTauriState>,
1080    connection_id: String,
1081) -> Result<Option<openrtc::client::ConnectionStateSnapshot>, String> {
1082    Ok(state.client().connection_state(&connection_id).await)
1083}
1084
1085#[tauri::command]
1086async fn start_rtc_device_subscription<R: Runtime>(
1087    state: tauri::State<'_, OpenRtcTauriState>,
1088    app: tauri::AppHandle<R>,
1089    request_id: String,
1090    user_id: String,
1091) -> Result<(), String> {
1092    stop_subscription_by_id(&state, &request_id).await;
1093    let mut stream = state
1094        .client()
1095        .subscribe_devices(&user_id)
1096        .await
1097        .map_err(|error| error.to_string())?;
1098    let request_for_task = request_id.clone();
1099    let handle = tokio::spawn(async move {
1100        while let Some(result) = stream.next().await {
1101            let payload = match result {
1102                Ok(events) => serde_json::json!({
1103                    "requestId": request_for_task,
1104                    "type": "deviceEvents",
1105                    "events": events,
1106                }),
1107                Err(error) => serde_json::json!({
1108                    "requestId": request_for_task,
1109                    "type": "error",
1110                    "message": error.to_string(),
1111                }),
1112            };
1113            let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
1114            let _ = app.emit(RTC_DEVICE_EVENTS, payload);
1115            if is_error {
1116                break;
1117            }
1118        }
1119    });
1120    state.subscriptions.lock().await.insert(request_id, handle);
1121    Ok(())
1122}
1123
1124#[tauri::command]
1125async fn start_rtc_session_subscription<R: Runtime>(
1126    state: tauri::State<'_, OpenRtcTauriState>,
1127    app: tauri::AppHandle<R>,
1128    request_id: String,
1129    local_device_id: String,
1130) -> Result<(), String> {
1131    stop_subscription_by_id(&state, &request_id).await;
1132    let mut stream = state
1133        .client()
1134        .subscribe_sessions(&local_device_id)
1135        .await
1136        .map_err(|error| error.to_string())?;
1137    let request_for_task = request_id.clone();
1138    let handle = tokio::spawn(async move {
1139        while let Some(result) = stream.next().await {
1140            let payload = match result {
1141                Ok(events) => serde_json::json!({
1142                    "requestId": request_for_task,
1143                    "type": "sessionEvents",
1144                    "events": events,
1145                }),
1146                Err(error) => serde_json::json!({
1147                    "requestId": request_for_task,
1148                    "type": "error",
1149                    "message": error.to_string(),
1150                }),
1151            };
1152            let is_error = payload.get("type").and_then(|value| value.as_str()) == Some("error");
1153            let _ = app.emit(RTC_SESSION_EVENTS, payload);
1154            if is_error {
1155                break;
1156            }
1157        }
1158    });
1159    state.subscriptions.lock().await.insert(request_id, handle);
1160    Ok(())
1161}
1162
1163#[tauri::command]
1164async fn stop_rtc_subscription(
1165    state: tauri::State<'_, OpenRtcTauriState>,
1166    request_id: String,
1167) -> Result<(), String> {
1168    stop_subscription_by_id(&state, &request_id).await;
1169    Ok(())
1170}
1171
1172async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
1173    if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
1174        handle.abort();
1175    }
1176}
1177
1178#[tauri::command]
1179async fn start_incoming_peer_bi_streams<R: Runtime>(
1180    app: tauri::AppHandle<R>,
1181    state: tauri::State<'_, OpenRtcTauriState>,
1182    request_id: Option<String>,
1183) -> Result<String, String> {
1184    let request_id = request_id
1185        .map(|value| value.trim().to_string())
1186        .filter(|value| !value.is_empty())
1187        .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1188
1189    stop_subscription_by_id(&state, &request_id).await;
1190
1191    let incoming = state
1192        .client()
1193        .incoming_streams()
1194        .await
1195        .map_err(|error| format!("failed to subscribe to incoming peer streams: {error}"))?;
1196    let app_for_task = app.clone();
1197    let request_for_task = request_id.clone();
1198    let handle = tokio::spawn(async move {
1199        while let Ok(incoming_stream) = incoming.recv().await {
1200            let remote_node_id = incoming_stream.endpoint_id.to_string();
1201            let IncomingStreamType::Bi(send, recv) = incoming_stream.stream else {
1202                continue;
1203            };
1204            let state_for_task = app_for_task.state::<OpenRtcTauriState>();
1205            let result = register_peer_bi_stream(
1206                app_for_task.clone(),
1207                state_for_task.inner(),
1208                None,
1209                remote_node_id.clone(),
1210                PeerSendStream::plain(send),
1211                PeerRecvStream::plain(recv),
1212                false,
1213            )
1214            .await;
1215            let _ = app_for_task.emit(
1216                PEER_BI_STREAM_INCOMING_EVENT,
1217                IncomingPeerBiStreamEvent {
1218                    request_id: request_for_task.clone(),
1219                    stream_id: result.stream_id,
1220                    remote_node_id,
1221                },
1222            );
1223        }
1224    });
1225
1226    state
1227        .subscriptions
1228        .lock()
1229        .await
1230        .insert(request_id.clone(), handle);
1231    Ok(request_id)
1232}
1233
1234#[tauri::command]
1235async fn open_peer_bi_stream<R: Runtime>(
1236    app: tauri::AppHandle<R>,
1237    state: tauri::State<'_, OpenRtcTauriState>,
1238    peer_id: String,
1239    timeout_ms: Option<u64>,
1240) -> Result<OpenPeerBiStreamResult, String> {
1241    open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
1242        client.open_peer_bi(&peer_id, timeout_ms).await
1243    })
1244    .await
1245}
1246
1247#[tauri::command]
1248async fn open_peer_bi_transport_only_stream<R: Runtime>(
1249    app: tauri::AppHandle<R>,
1250    state: tauri::State<'_, OpenRtcTauriState>,
1251    peer_id: String,
1252    timeout_ms: Option<u64>,
1253) -> Result<OpenPeerBiStreamResult, String> {
1254    open_peer_bi_with(
1255        app,
1256        state,
1257        peer_id,
1258        timeout_ms,
1259        |client, peer_id, timeout_ms| async move {
1260            client
1261                .open_peer_bi_transport_only(&peer_id, timeout_ms)
1262                .await
1263                .map(|(connection_id, remote_node_id, send, recv)| {
1264                    (
1265                        connection_id,
1266                        remote_node_id,
1267                        PeerSendStream::plain(send),
1268                        PeerRecvStream::plain(recv),
1269                    )
1270                })
1271        },
1272    )
1273    .await
1274}
1275
1276#[tauri::command]
1277async fn open_peer_native_bi_stream<R: Runtime>(
1278    app: tauri::AppHandle<R>,
1279    state: tauri::State<'_, OpenRtcTauriState>,
1280    peer_id: String,
1281    label: String,
1282    timeout_ms: Option<u64>,
1283) -> Result<OpenPeerBiStreamResult, String> {
1284    let channel_envelope = openrtc::stream_metadata::encode_channel_envelope(&label, None)
1285        .map_err(|error| error.to_string())?;
1286    open_peer_bi_with(
1287        app,
1288        state,
1289        peer_id,
1290        timeout_ms,
1291        move |client, peer_id, timeout_ms| async move {
1292            let (connection_id, remote_node_id, mut send, recv) = client
1293                .open_peer_bi_transport_only(&peer_id, timeout_ms)
1294                .await?;
1295            send.write_all(&channel_envelope).await?;
1296            Ok((
1297                connection_id,
1298                remote_node_id,
1299                PeerSendStream::plain(send),
1300                PeerRecvStream::plain(recv),
1301            ))
1302        },
1303    )
1304    .await
1305}
1306
1307#[tauri::command]
1308async fn write_peer_bi_stream(
1309    state: tauri::State<'_, OpenRtcTauriState>,
1310    stream_id: String,
1311    bytes: Vec<u8>,
1312) -> Result<(), String> {
1313    let stream_id = stream_id.trim().to_string();
1314    if stream_id.is_empty() {
1315        return Err("streamId is required".to_string());
1316    }
1317
1318    let send = {
1319        let streams = state.peer_bi_streams.lock().await;
1320        streams
1321            .get(&stream_id)
1322            .map(|handle| handle.send.clone())
1323            .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
1324    };
1325    let result = send
1326        .lock()
1327        .await
1328        .as_mut()
1329        .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
1330        .write_all(&bytes)
1331        .await
1332        .map_err(|error| format!("write peer bi stream failed: {error}"));
1333    result
1334}
1335
1336#[tauri::command]
1337async fn start_peer_bi_stream_read<R: Runtime>(
1338    app: tauri::AppHandle<R>,
1339    state: tauri::State<'_, OpenRtcTauriState>,
1340    stream_id: String,
1341) -> Result<(), String> {
1342    let stream_id = stream_id.trim().to_string();
1343    if stream_id.is_empty() {
1344        return Err("streamId is required".to_string());
1345    }
1346
1347    let mut streams = state.peer_bi_streams.lock().await;
1348    let handle = streams
1349        .get_mut(&stream_id)
1350        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
1351    if handle.read_task.is_some() {
1352        return Ok(());
1353    }
1354    let recv = handle
1355        .recv
1356        .take()
1357        .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
1358    handle.read_task = Some(spawn_peer_bi_stream_reader(app, stream_id, recv));
1359    Ok(())
1360}
1361
1362#[tauri::command]
1363async fn close_peer_bi_stream(
1364    state: tauri::State<'_, OpenRtcTauriState>,
1365    stream_id: String,
1366) -> Result<(), String> {
1367    let stream_id = stream_id.trim().to_string();
1368    if stream_id.is_empty() {
1369        return Err("streamId is required".to_string());
1370    }
1371
1372    let handle = {
1373        let mut streams = state.peer_bi_streams.lock().await;
1374        streams.remove(&stream_id)
1375    };
1376    if let Some(handle) = handle {
1377        // Hold the writer mutex before finishing so close waits for any
1378        // in-flight chunk write instead of skipping the flush when Arc clones exist.
1379        let mut send = handle.send.lock().await;
1380        if let Some(send) = send.take() {
1381            let _ = send.finish();
1382        }
1383        drop(send);
1384        if let Some(read_task) = handle.read_task {
1385            read_task.abort();
1386        }
1387    }
1388    Ok(())
1389}
1390
1391#[tauri::command]
1392async fn get_app_limits() -> Result<serde_json::Value, String> {
1393    Ok(serde_json::json!({
1394        "devicesPerUser": -1,
1395        "maxRooms": -1,
1396        "maxMembersPerRoom": -1
1397    }))
1398}
1399
1400#[cfg(test)]
1401mod tests {
1402    use super::*;
1403
1404    #[test]
1405    fn derives_app_tag_from_api_key_by_default() {
1406        let config = OpenRtcTauriConfig {
1407            api_key: Some("pk_test_1234567890abcdef".to_string()),
1408            ..OpenRtcTauriConfig::default()
1409        };
1410
1411        assert_eq!(config.app_tag(), "app_1234567890abcdef");
1412    }
1413
1414    #[test]
1415    fn explicit_app_tag_wins_over_api_key() {
1416        let config = OpenRtcTauriConfig {
1417            api_key: Some("pk_test_1234567890abcdef".to_string()),
1418            app_tag: Some("space::manual".to_string()),
1419            ..OpenRtcTauriConfig::default()
1420        };
1421
1422        assert_eq!(config.app_tag(), "space::manual");
1423    }
1424
1425    #[test]
1426    fn derives_space_app_tag_from_space_key() {
1427        let config = OpenRtcTauriConfig {
1428            api_key: Some("pk_test_abc".to_string()),
1429            space_key: Some("room".to_string()),
1430            ..OpenRtcTauriConfig::default()
1431        };
1432
1433        assert!(config.app_tag().starts_with("space::"));
1434        assert_eq!(config.app_tag().len(), "space::".len() + 64);
1435    }
1436
1437    #[test]
1438    fn native_state_reconfigures_anonymous_client_to_public_space() {
1439        let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
1440        assert_eq!(state.client().app_tag(), "app_anonymous");
1441
1442        let (client, reconfigured) = state.configure_identity(
1443            Some("pk_test_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string()),
1444            Some("assistant-control-plane".to_string()),
1445            None,
1446        );
1447
1448        assert!(reconfigured);
1449        assert!(client.app_tag().starts_with("space::"));
1450        assert_eq!(client.app_tag(), state.client().app_tag());
1451        assert_ne!(client.app_tag(), "app_anonymous");
1452    }
1453
1454    #[test]
1455    fn managed_session_prefers_requested_logical_device_id() {
1456        let identity = openrtc::native_device::NativeDeviceIdentity {
1457            device_id: "persisted-native-device".to_string(),
1458            device_name: "Mac".to_string(),
1459            created_at_ms: 1,
1460            updated_at_ms: 1,
1461            name_source: None,
1462            system_info: None,
1463        };
1464
1465        assert_eq!(
1466            managed_session_device_id(Some(" desktop-e2e-native "), &identity),
1467            "desktop-e2e-native"
1468        );
1469        assert_eq!(
1470            managed_session_device_id(None, &identity),
1471            "persisted-native-device"
1472        );
1473    }
1474
1475    #[test]
1476    fn managed_session_metadata_uses_authoritative_device_id() {
1477        let raw = serde_json::json!({
1478            "deviceId": "persisted-native-device",
1479            "assistantDevice": {
1480                "deviceId": "desktop-e2e-native"
1481            }
1482        })
1483        .to_string();
1484
1485        let metadata =
1486            metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
1487        let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
1488
1489        assert_eq!(parsed["deviceId"], "desktop-e2e-native");
1490        assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
1491    }
1492}