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