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