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