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