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