Skip to main content

openrtc_tauri_plugin/
lib.rs

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