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