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 remote_node_id = incoming_stream.endpoint_id.to_string();
1763            let transport_stable_id = incoming_stream.transport_stable_id;
1764            let IncomingStreamType::Bi(send, recv) = incoming_stream.stream else {
1765                continue;
1766            };
1767            if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1768                eprintln!(
1769                    "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} request_id={} phase=received",
1770                    remote_node_id,
1771                    transport_stable_id,
1772                    request_for_task
1773                );
1774            }
1775            let state_for_task = app_for_task.state::<OpenRtcTauriState>();
1776            let (send, recv) = match state_for_task
1777                .client()
1778                .wrap_incoming_application_bi_stream_for_transport(
1779                    &incoming_stream.endpoint_id,
1780                    transport_stable_id,
1781                    send,
1782                    recv,
1783                )
1784                .await
1785            {
1786                Ok(streams) => streams,
1787                Err(error) => {
1788                    eprintln!(
1789                        "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} phase=rejected error={}",
1790                        remote_node_id, transport_stable_id, error
1791                    );
1792                    continue;
1793                }
1794            };
1795            let result = register_peer_bi_stream(
1796                app_for_task.clone(),
1797                state_for_task.inner(),
1798                None,
1799                remote_node_id.clone(),
1800                send,
1801                recv,
1802                false,
1803            )
1804            .await;
1805            let _ = app_for_task.emit(
1806                PEER_BI_STREAM_INCOMING_EVENT,
1807                IncomingPeerBiStreamEvent {
1808                    request_id: request_for_task.clone(),
1809                    stream_id: result.stream_id,
1810                    remote_node_id: remote_node_id.clone(),
1811                    transport_stable_id,
1812                },
1813            );
1814            if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1815                eprintln!(
1816                    "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} request_id={} phase=emitted",
1817                    remote_node_id,
1818                    transport_stable_id,
1819                    request_for_task
1820                );
1821            }
1822        }
1823    });
1824
1825    state
1826        .subscriptions
1827        .lock()
1828        .await
1829        .insert(request_id.clone(), handle);
1830    Ok(request_id)
1831}
1832
1833#[tauri::command]
1834async fn is_current_transport_stable_id(
1835    state: tauri::State<'_, OpenRtcTauriState>,
1836    endpoint_id: String,
1837    transport_stable_id: u64,
1838) -> Result<bool, String> {
1839    let is_current = state
1840        .client()
1841        .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
1842        .await
1843        .map_err(|error| error.to_string())?;
1844    if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1845        eprintln!(
1846            "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
1847            endpoint_id, transport_stable_id, is_current
1848        );
1849    }
1850    Ok(is_current)
1851}
1852
1853#[tauri::command]
1854async fn open_peer_bi_stream<R: Runtime>(
1855    app: tauri::AppHandle<R>,
1856    state: tauri::State<'_, OpenRtcTauriState>,
1857    peer_id: String,
1858    timeout_ms: Option<u64>,
1859) -> Result<OpenPeerBiStreamResult, String> {
1860    open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
1861        client.open_peer_bi(&peer_id, timeout_ms).await
1862    })
1863    .await
1864}
1865
1866#[tauri::command]
1867async fn open_peer_bi_transport_only_stream<R: Runtime>(
1868    app: tauri::AppHandle<R>,
1869    state: tauri::State<'_, OpenRtcTauriState>,
1870    peer_id: String,
1871    timeout_ms: Option<u64>,
1872) -> Result<OpenPeerBiStreamResult, String> {
1873    open_peer_bi_with(
1874        app,
1875        state,
1876        peer_id,
1877        timeout_ms,
1878        |client, peer_id, timeout_ms| async move {
1879            client
1880                .open_peer_bi_transport_only(&peer_id, timeout_ms)
1881                .await
1882                .map(|(connection_id, remote_node_id, send, recv)| {
1883                    (
1884                        connection_id,
1885                        remote_node_id,
1886                        PeerSendStream::plain(send),
1887                        PeerRecvStream::plain(recv),
1888                    )
1889                })
1890        },
1891    )
1892    .await
1893}
1894
1895#[tauri::command]
1896async fn open_peer_native_bi_stream<R: Runtime>(
1897    app: tauri::AppHandle<R>,
1898    state: tauri::State<'_, OpenRtcTauriState>,
1899    peer_id: String,
1900    label: String,
1901    timeout_ms: Option<u64>,
1902) -> Result<OpenPeerBiStreamResult, String> {
1903    let channel_envelope = openrtc::stream_metadata::encode_channel_envelope(&label, None)
1904        .map_err(|error| error.to_string())?;
1905    open_peer_bi_with(
1906        app,
1907        state,
1908        peer_id,
1909        timeout_ms,
1910        move |client, peer_id, timeout_ms| async move {
1911            let (connection_id, remote_node_id, mut send, recv) =
1912                client.open_peer_bi(&peer_id, timeout_ms).await?;
1913            send.write_all(&channel_envelope).await?;
1914            Ok((connection_id, remote_node_id, send, recv))
1915        },
1916    )
1917    .await
1918}
1919
1920#[tauri::command]
1921async fn write_peer_bi_stream(
1922    state: tauri::State<'_, OpenRtcTauriState>,
1923    stream_id: String,
1924    bytes: Vec<u8>,
1925) -> Result<(), String> {
1926    let stream_id = stream_id.trim().to_string();
1927    if stream_id.is_empty() {
1928        return Err("streamId is required".to_string());
1929    }
1930
1931    let send = {
1932        let streams = state.peer_bi_streams.lock().await;
1933        streams
1934            .get(&stream_id)
1935            .map(|handle| handle.send.clone())
1936            .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
1937    };
1938    let result = send
1939        .lock()
1940        .await
1941        .as_mut()
1942        .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
1943        .write_all(&bytes)
1944        .await
1945        .map_err(|error| format!("write peer bi stream failed: {error}"));
1946    result
1947}
1948
1949#[tauri::command]
1950async fn start_peer_bi_stream_read<R: Runtime>(
1951    app: tauri::AppHandle<R>,
1952    state: tauri::State<'_, OpenRtcTauriState>,
1953    stream_id: String,
1954) -> Result<(), String> {
1955    let stream_id = stream_id.trim().to_string();
1956    if stream_id.is_empty() {
1957        return Err("streamId is required".to_string());
1958    }
1959
1960    let mut streams = state.peer_bi_streams.lock().await;
1961    let handle = streams
1962        .get_mut(&stream_id)
1963        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
1964    if handle.read_task.is_some() {
1965        return Ok(());
1966    }
1967    let recv = handle
1968        .recv
1969        .take()
1970        .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
1971    handle.read_task = Some(spawn_peer_bi_stream_reader(app, stream_id, recv));
1972    Ok(())
1973}
1974
1975#[tauri::command]
1976async fn close_peer_bi_stream(
1977    state: tauri::State<'_, OpenRtcTauriState>,
1978    stream_id: String,
1979) -> Result<(), String> {
1980    let stream_id = stream_id.trim().to_string();
1981    if stream_id.is_empty() {
1982        return Err("streamId is required".to_string());
1983    }
1984
1985    let handle = {
1986        let mut streams = state.peer_bi_streams.lock().await;
1987        streams.remove(&stream_id)
1988    };
1989    if let Some(handle) = handle {
1990        // Hold the writer mutex before finishing so close waits for any
1991        // in-flight chunk write instead of skipping the flush when Arc clones exist.
1992        let mut send = handle.send.lock().await;
1993        if let Some(send) = send.take() {
1994            let _ = send.finish();
1995        }
1996        drop(send);
1997        if let Some(read_task) = handle.read_task {
1998            read_task.abort();
1999        }
2000    }
2001    Ok(())
2002}
2003
2004#[tauri::command]
2005async fn get_app_limits() -> Result<serde_json::Value, String> {
2006    Ok(serde_json::json!({
2007        "devicesPerUser": -1,
2008        "maxRooms": -1,
2009        "maxMembersPerRoom": -1
2010    }))
2011}
2012
2013#[cfg(test)]
2014mod tests {
2015    use super::*;
2016    use std::sync::atomic::{AtomicUsize, Ordering};
2017    use tokio::sync::oneshot;
2018
2019    struct FakeNativeTransportInstaller {
2020        installs: Arc<AtomicUsize>,
2021    }
2022
2023    impl NativeTransportInstaller for FakeNativeTransportInstaller {
2024        fn id(&self) -> &'static str {
2025            "ble"
2026        }
2027
2028        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2029            config.ble.as_ref().is_some_and(|ble| ble.enabled)
2030        }
2031
2032        fn install(
2033            &self,
2034            _client: Arc<openrtc::client::Client>,
2035            _config: openrtc::client::TransportConfig,
2036            _context: NativeTransportInstallContext,
2037        ) -> NativeTransportInstallFuture {
2038            self.installs.fetch_add(1, Ordering::SeqCst);
2039            Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
2040        }
2041    }
2042
2043    #[derive(Debug)]
2044    struct UnavailableNativeTransportInstaller;
2045
2046    impl NativeTransportInstaller for UnavailableNativeTransportInstaller {
2047        fn id(&self) -> &'static str {
2048            "ble"
2049        }
2050
2051        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2052            config.ble.as_ref().is_some_and(|ble| ble.enabled)
2053        }
2054
2055        fn install(
2056            &self,
2057            _client: Arc<openrtc::client::Client>,
2058            _config: openrtc::client::TransportConfig,
2059            _context: NativeTransportInstallContext,
2060        ) -> NativeTransportInstallFuture {
2061            Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
2062        }
2063    }
2064
2065    fn requested_ble_config() -> openrtc::client::TransportConfig {
2066        openrtc::client::TransportConfig {
2067            ble: Some(openrtc::client::BleConfig {
2068                enabled: true,
2069                ..openrtc::client::BleConfig::default()
2070            }),
2071            ..openrtc::client::TransportConfig::default()
2072        }
2073    }
2074
2075    #[test]
2076    fn native_trust_relay_stays_fail_closed_after_clear() {
2077        let relay = Arc::new(TokenRelayState::default());
2078        let provider = relay.native_trust_token_provider();
2079        assert!(provider().is_none());
2080
2081        relay.set_native_trust_token(
2082            Some("opaque-trust-token".to_string()),
2083            Some(1_900_000_000_000),
2084        );
2085        assert_eq!(
2086            provider(),
2087            Some(openrtc::client::NativeTrustToken {
2088                token: "opaque-trust-token".to_string(),
2089                expires_at_ms: 1_900_000_000_000,
2090            }),
2091        );
2092
2093        relay.set_native_trust_token(None, None);
2094        assert_eq!(
2095            provider(),
2096            Some(openrtc::client::NativeTrustToken {
2097                token: String::new(),
2098                expires_at_ms: 0,
2099            }),
2100        );
2101    }
2102
2103    fn native_transport_install_context() -> NativeTransportInstallContext {
2104        NativeTransportInstallContext {
2105            data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
2106        }
2107    }
2108
2109    fn managed_session_test_result(user_id: &str) -> StartRtcManagedSessionResult {
2110        StartRtcManagedSessionResult {
2111            local_node_id: format!("node-{user_id}"),
2112            ticket_scope: Some("user-device".to_string()),
2113            ticket: None,
2114            presence_started: true,
2115            auto_connect_started: true,
2116            local_device: openrtc::native_device::NativeDeviceIdentity {
2117                device_id: format!("device-{user_id}"),
2118                device_name: "Test Device".to_string(),
2119                created_at_ms: 1,
2120                updated_at_ms: 1,
2121                name_source: None,
2122                system_info: None,
2123            },
2124        }
2125    }
2126
2127    async fn simulate_managed_session_start(
2128        state: Arc<OpenRtcTauriState>,
2129        api_key: String,
2130        user_id: String,
2131        pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
2132        starts: Arc<AtomicUsize>,
2133        stopped_owners: Arc<Mutex<Vec<usize>>>,
2134    ) -> ManagedSessionDisposition {
2135        let _start_guard = state.managed_session_start_guard.lock().await;
2136        let (client, _) = state.configure_identity(Some(api_key), None, None);
2137        let key = format!("{}:{user_id}:device-{user_id}", client.app_tag());
2138        let (disposition, active_epoch) = {
2139            let active = state.managed_session.lock().await;
2140            let disposition = managed_session_disposition(
2141                active.as_ref().map(|record| record.key.as_str()),
2142                active
2143                    .as_ref()
2144                    .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
2145                active
2146                    .as_ref()
2147                    .is_some_and(|record| record.result.presence_started),
2148                active
2149                    .as_ref()
2150                    .is_some_and(|record| record.result.auto_connect_started),
2151                &key,
2152                true,
2153                true,
2154            );
2155            (
2156                disposition,
2157                active.as_ref().map(|record| record.owner_epoch),
2158            )
2159        };
2160
2161        if disposition == ManagedSessionDisposition::Reuse {
2162            return disposition;
2163        }
2164
2165        if disposition == ManagedSessionDisposition::Replace {
2166            let previous = state
2167                .managed_session
2168                .lock()
2169                .await
2170                .take()
2171                .expect("replacement must have a previous owner");
2172            stopped_owners
2173                .lock()
2174                .await
2175                .push(Arc::as_ptr(&previous.owner_client) as usize);
2176        }
2177
2178        starts.fetch_add(1, Ordering::SeqCst);
2179        if let Some((entered, release)) = pause {
2180            entered.send(()).expect("start observer must be waiting");
2181            release.await.expect("start release must be sent");
2182        }
2183
2184        let owner_epoch = match disposition {
2185            ManagedSessionDisposition::Refresh => {
2186                active_epoch.expect("refresh must preserve the active owner epoch")
2187            }
2188            ManagedSessionDisposition::Start | ManagedSessionDisposition::Replace => {
2189                state.allocate_managed_session_owner_epoch()
2190            }
2191            ManagedSessionDisposition::Reuse => unreachable!("reuse returned before startup"),
2192        };
2193        let result = managed_session_test_result(&user_id);
2194        *state.managed_session.lock().await = Some(ManagedSessionRecord {
2195            key,
2196            owner_client: client,
2197            owner_epoch,
2198            result,
2199        });
2200        disposition
2201    }
2202
2203    #[tokio::test]
2204    async fn concurrent_same_key_start_reuses_one_owner_epoch() {
2205        let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
2206        let starts = Arc::new(AtomicUsize::new(0));
2207        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
2208        let api_key = "pk_test_same_managed_session_owner".to_string();
2209        let (first_entered_tx, first_entered_rx) = oneshot::channel();
2210        let (first_release_tx, first_release_rx) = oneshot::channel();
2211
2212        let first = tokio::spawn(simulate_managed_session_start(
2213            state.clone(),
2214            api_key.clone(),
2215            "same-user".to_string(),
2216            Some((first_entered_tx, first_release_rx)),
2217            starts.clone(),
2218            stopped_owners.clone(),
2219        ));
2220        first_entered_rx.await.expect("first start must pause");
2221
2222        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
2223        let second_state = state.clone();
2224        let second_starts = starts.clone();
2225        let second_stopped_owners = stopped_owners.clone();
2226        let second = tokio::spawn(async move {
2227            second_attempting_tx.send(()).unwrap();
2228            simulate_managed_session_start(
2229                second_state,
2230                api_key,
2231                "same-user".to_string(),
2232                None,
2233                second_starts,
2234                second_stopped_owners,
2235            )
2236            .await
2237        });
2238        second_attempting_rx.await.unwrap();
2239        tokio::task::yield_now().await;
2240        assert_eq!(starts.load(Ordering::SeqCst), 1);
2241
2242        first_release_tx.send(()).unwrap();
2243        assert_eq!(
2244            first.await.expect("first start task"),
2245            ManagedSessionDisposition::Start
2246        );
2247        assert_eq!(
2248            second.await.expect("second start task"),
2249            ManagedSessionDisposition::Reuse
2250        );
2251        assert_eq!(starts.load(Ordering::SeqCst), 1);
2252
2253        let active = state
2254            .managed_session
2255            .lock()
2256            .await
2257            .clone()
2258            .expect("active owner");
2259        assert_eq!(active.owner_epoch, 1);
2260        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
2261        assert!(stopped_owners.lock().await.is_empty());
2262    }
2263
2264    #[tokio::test]
2265    async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
2266        let state = Arc::new(OpenRtcTauriState::new(OpenRtcTauriConfig::default()));
2267        let starts = Arc::new(AtomicUsize::new(0));
2268        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
2269        let (first_entered_tx, first_entered_rx) = oneshot::channel();
2270        let (first_release_tx, first_release_rx) = oneshot::channel();
2271
2272        let first = tokio::spawn(simulate_managed_session_start(
2273            state.clone(),
2274            "pk_test_owner_a".to_string(),
2275            "first-user".to_string(),
2276            Some((first_entered_tx, first_release_rx)),
2277            starts.clone(),
2278            stopped_owners.clone(),
2279        ));
2280        first_entered_rx.await.expect("first start must pause");
2281        let first_owner = state.client();
2282
2283        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
2284        let second_state = state.clone();
2285        let second_starts = starts.clone();
2286        let second_stopped_owners = stopped_owners.clone();
2287        let second = tokio::spawn(async move {
2288            second_attempting_tx.send(()).unwrap();
2289            simulate_managed_session_start(
2290                second_state,
2291                "pk_test_owner_b".to_string(),
2292                "second-user".to_string(),
2293                None,
2294                second_starts,
2295                second_stopped_owners,
2296            )
2297            .await
2298        });
2299        second_attempting_rx.await.unwrap();
2300        tokio::task::yield_now().await;
2301        assert!(Arc::ptr_eq(&first_owner, &state.client()));
2302
2303        first_release_tx.send(()).unwrap();
2304        assert_eq!(
2305            first.await.expect("first start task"),
2306            ManagedSessionDisposition::Start
2307        );
2308        assert_eq!(
2309            second.await.expect("replacement start task"),
2310            ManagedSessionDisposition::Replace
2311        );
2312        assert_eq!(starts.load(Ordering::SeqCst), 2);
2313
2314        let active = state
2315            .managed_session
2316            .lock()
2317            .await
2318            .clone()
2319            .expect("active owner");
2320        assert_eq!(active.owner_epoch, 2);
2321        assert!(active.key.contains("second-user"));
2322        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
2323        assert_eq!(
2324            stopped_owners.lock().await.as_slice(),
2325            [Arc::as_ptr(&first_owner) as usize]
2326        );
2327        assert!(!Arc::ptr_eq(&active.owner_client, &first_owner));
2328    }
2329
2330    #[test]
2331    fn derives_app_tag_from_api_key_by_default() {
2332        let config = OpenRtcTauriConfig {
2333            api_key: Some("pk_test_1234567890abcdef".to_string()),
2334            ..OpenRtcTauriConfig::default()
2335        };
2336
2337        assert_eq!(config.app_tag(), "app_1234567890abcdef");
2338    }
2339
2340    #[tokio::test]
2341    async fn requested_native_transport_without_installer_keeps_base_route_available() {
2342        let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
2343        let client = state.client();
2344        ensure_requested_native_transports(
2345            &state,
2346            &client,
2347            Some(&requested_ble_config()),
2348            &native_transport_install_context(),
2349        )
2350        .await
2351        .expect("an unavailable optional transport must not fail base Iroh startup");
2352
2353        assert!(state.installed_native_transports.lock().await.is_empty());
2354    }
2355
2356    #[tokio::test]
2357    async fn native_transport_installer_is_idempotent_for_one_client() {
2358        let installs = Arc::new(AtomicUsize::new(0));
2359        let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
2360            .with_native_transport_installer(Arc::new(FakeNativeTransportInstaller {
2361                installs: installs.clone(),
2362            }));
2363        let client = state.client();
2364        let config = requested_ble_config();
2365
2366        ensure_requested_native_transports(
2367            &state,
2368            &client,
2369            Some(&config),
2370            &native_transport_install_context(),
2371        )
2372        .await
2373        .expect("first install");
2374        ensure_requested_native_transports(
2375            &state,
2376            &client,
2377            Some(&config),
2378            &native_transport_install_context(),
2379        )
2380        .await
2381        .expect("idempotent install");
2382
2383        assert_eq!(installs.load(Ordering::SeqCst), 1);
2384    }
2385
2386    #[tokio::test]
2387    async fn unavailable_native_transport_keeps_base_route_available() {
2388        let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default())
2389            .with_native_transport_installer(Arc::new(UnavailableNativeTransportInstaller));
2390        let client = state.client();
2391
2392        ensure_requested_native_transports(
2393            &state,
2394            &client,
2395            Some(&requested_ble_config()),
2396            &native_transport_install_context(),
2397        )
2398        .await
2399        .expect("optional transport installation failure must not fail base Iroh startup");
2400
2401        assert!(state.installed_native_transports.lock().await.is_empty());
2402    }
2403
2404    #[test]
2405    fn explicit_app_tag_wins_over_api_key() {
2406        let config = OpenRtcTauriConfig {
2407            api_key: Some("pk_test_1234567890abcdef".to_string()),
2408            app_tag: Some("space::manual".to_string()),
2409            ..OpenRtcTauriConfig::default()
2410        };
2411
2412        assert_eq!(config.app_tag(), "space::manual");
2413    }
2414
2415    #[test]
2416    fn derives_space_app_tag_from_space_key() {
2417        let config = OpenRtcTauriConfig {
2418            api_key: Some("pk_test_abc".to_string()),
2419            space_key: Some("room".to_string()),
2420            ..OpenRtcTauriConfig::default()
2421        };
2422
2423        assert!(config.app_tag().starts_with("space::"));
2424        assert_eq!(config.app_tag().len(), "space::".len() + 64);
2425    }
2426
2427    #[test]
2428    fn from_env_accepts_vite_project_and_app_tag() {
2429        let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
2430        let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
2431        std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
2432        std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
2433
2434        let config = OpenRtcTauriConfig::from_env();
2435
2436        assert_eq!(config.project_id, "pluto-rtc-prod");
2437        assert_eq!(config.app_tag(), "app_from_vite_env");
2438
2439        match previous_project {
2440            Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
2441            None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
2442        }
2443        match previous_app_tag {
2444            Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
2445            None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
2446        }
2447    }
2448
2449    #[test]
2450    fn native_state_reconfigures_anonymous_client_to_public_space() {
2451        let state = OpenRtcTauriState::new(OpenRtcTauriConfig::default());
2452        assert_eq!(state.client().app_tag(), "app_anonymous");
2453
2454        let (client, reconfigured) = state.configure_identity(
2455            Some("pk_test_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string()),
2456            Some("assistant-control-plane".to_string()),
2457            None,
2458        );
2459
2460        assert!(reconfigured);
2461        assert!(client.app_tag().starts_with("space::"));
2462        assert_eq!(client.app_tag(), state.client().app_tag());
2463        assert_ne!(client.app_tag(), "app_anonymous");
2464    }
2465
2466    #[test]
2467    fn managed_session_uses_persisted_native_device_id() {
2468        let identity = openrtc::native_device::NativeDeviceIdentity {
2469            device_id: "persisted-native-device".to_string(),
2470            device_name: "Mac".to_string(),
2471            created_at_ms: 1,
2472            updated_at_ms: 1,
2473            name_source: None,
2474            system_info: None,
2475        };
2476
2477        assert_eq!(
2478            managed_session_device_id(Some(" desktop-e2e-native "), &identity),
2479            "persisted-native-device"
2480        );
2481        assert_eq!(
2482            managed_session_device_id(None, &identity),
2483            "persisted-native-device"
2484        );
2485    }
2486
2487    #[test]
2488    fn repeated_managed_session_triggers_converge_on_one_owner() {
2489        let key = "app:user:persisted-native-device";
2490        let cases = [
2491            ("react-remount", true, true, true, true),
2492            ("hmr", true, true, true, true),
2493            ("auth-refresh", true, true, true, true),
2494            ("resume", true, true, true, true),
2495            ("alias-change", true, true, true, true),
2496        ];
2497        for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
2498            assert_eq!(
2499                managed_session_disposition(
2500                    Some(key),
2501                    true,
2502                    active_presence,
2503                    active_auto,
2504                    key,
2505                    requested_presence,
2506                    requested_auto,
2507                ),
2508                ManagedSessionDisposition::Reuse,
2509                "{label} must reuse the authoritative tuple"
2510            );
2511        }
2512
2513        assert_eq!(
2514            managed_session_disposition(Some(key), true, false, true, key, true, true),
2515            ManagedSessionDisposition::Refresh,
2516            "failed presence startup must retry idempotently"
2517        );
2518        assert_eq!(
2519            managed_session_disposition(
2520                Some(key),
2521                true,
2522                true,
2523                true,
2524                "app:other-user:persisted-native-device",
2525                true,
2526                true,
2527            ),
2528            ManagedSessionDisposition::Replace,
2529            "user switching must replace the previous lifecycle owner"
2530        );
2531        assert_eq!(
2532            managed_session_disposition(Some(key), false, true, true, key, true, true),
2533            ManagedSessionDisposition::Replace,
2534            "the same tuple on a different client epoch must replace the previous owner"
2535        );
2536    }
2537
2538    #[test]
2539    fn managed_session_metadata_uses_authoritative_device_id() {
2540        let raw = serde_json::json!({
2541            "deviceId": "persisted-native-device",
2542            "assistantDevice": {
2543                "deviceId": "desktop-e2e-native"
2544            }
2545        })
2546        .to_string();
2547
2548        let metadata =
2549            metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
2550        let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
2551
2552        assert_eq!(parsed["deviceId"], "desktop-e2e-native");
2553        assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
2554    }
2555}