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