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