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