Skip to main content

openrtc_tauri_plugin/
lib.rs

1use std::collections::{BTreeSet, 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 serde::{Deserialize, Serialize};
9use tauri::ipc::{Channel, InvokeBody, Request, Response};
10use tauri::{Manager, Runtime};
11use tokio::sync::Mutex;
12
13const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
14const PEER_BI_STREAM_ID_HEADER: &str = "x-openrtc-stream-id";
15const PEER_ID_HEADER: &str = "x-openrtc-peer-id";
16const MEDIA_PUBLICATION_ID_HEADER: &str = "x-openrtc-publication-id";
17const MEDIA_TIMESTAMP_US_HEADER: &str = "x-openrtc-timestamp-us";
18const MEDIA_DURATION_US_HEADER: &str = "x-openrtc-duration-us";
19const MEDIA_KEYFRAME_HEADER: &str = "x-openrtc-keyframe";
20const MEDIA_DISCARDABLE_HEADER: &str = "x-openrtc-discardable";
21
22fn request_body_bytes(request: &Request<'_>) -> Result<Vec<u8>, String> {
23    match request.body() {
24        InvokeBody::Raw(bytes) => Ok(bytes.clone()),
25        // Tauri Android can forward typed-array command bodies through JSON.
26        // Preserve one command contract without falling back to global events.
27        InvokeBody::Json(json) => serde_json::from_value::<Vec<u8>>(json.clone())
28            .map_err(|error| format!("invalid binary IPC payload: {error}")),
29    }
30}
31
32fn required_request_header(request: &Request<'_>, name: &str) -> Result<String, String> {
33    request
34        .headers()
35        .get(name)
36        .and_then(|value| value.to_str().ok())
37        .map(str::trim)
38        .filter(|value| !value.is_empty())
39        .map(ToOwned::to_owned)
40        .ok_or_else(|| format!("{name} header is required"))
41}
42
43fn parsed_request_header<T>(request: &Request<'_>, name: &str) -> Result<T, String>
44where
45    T: std::str::FromStr,
46    T::Err: std::fmt::Display,
47{
48    required_request_header(request, name)?
49        .parse::<T>()
50        .map_err(|error| format!("invalid {name} header: {error}"))
51}
52
53pub type InstallFuture = std::pin::Pin<
54    Box<
55        dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
56            + Send,
57    >,
58>;
59
60#[derive(Debug, Clone)]
61pub struct InstallContext {
62    /// Stable host-owned storage used by every native endpoint initializer.
63    /// Installers must derive their Iroh identity from this directory rather
64    /// than creating a transport-specific identity.
65    pub data_dir: PathBuf,
66}
67
68/// Host-supplied implementation for a native transport that cannot be bundled
69/// in the publishable Tauri plugin. The OpenRTC constructor remains the sole
70/// activation switch; registering an installer only declares compiled support.
71pub trait TransportInstaller: Send + Sync {
72    fn id(&self) -> &'static str;
73    fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
74    fn install(
75        &self,
76        client: Arc<openrtc::client::Client>,
77        config: openrtc::client::TransportConfig,
78        context: InstallContext,
79    ) -> InstallFuture;
80}
81
82/// Host-owned private-storage bridge for the OpenRTC 2.0 per-install device
83/// proof key. OpenRTC never receives private key bytes and does not prescribe
84/// an interactive OS credential store. The host exposes only public JWK export
85/// and signing; Plutonium's current shipping owner is its prompt-free,
86/// app-private store.
87pub trait DeviceKeySigner: Send + Sync {
88    fn public_jwk(&self, app_tag: &str) -> Result<serde_json::Value, String>;
89    fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String>;
90    fn delete(&self, app_tag: &str) -> Result<(), String>;
91    fn read_secure_record(&self, _app_tag: &str, _key: &str) -> Result<Option<String>, String> {
92        Ok(None)
93    }
94    fn write_secure_record(&self, _app_tag: &str, _key: &str, _value: &str) -> Result<(), String> {
95        Err("OpenRTC 2.0 native certificate persistence requires a host secure store".to_string())
96    }
97    fn delete_secure_record(&self, _app_tag: &str, _key: &str) -> Result<(), String> {
98        Ok(())
99    }
100}
101
102struct InstalledNativeTransport {
103    client_ptr: usize,
104    _runtime: Box<dyn std::any::Any + Send + Sync>,
105}
106
107#[derive(Debug, Clone, Serialize, Deserialize)]
108#[serde(rename_all = "camelCase")]
109pub struct OpenRtcTauriConfig {
110    /// Public OpenRTC developer API key. The app identity is derived locally;
111    /// consumers cannot retarget the native runtime with an app tag or
112    /// provider project identifier.
113    pub api_key: String,
114    #[serde(default)]
115    pub data_dir: Option<PathBuf>,
116    #[serde(default)]
117    pub transport_config: Option<openrtc::client::TransportConfig>,
118}
119
120impl Default for OpenRtcTauriConfig {
121    fn default() -> Self {
122        Self {
123            api_key: String::new(),
124            data_dir: None,
125            transport_config: None,
126        }
127    }
128}
129
130impl OpenRtcTauriConfig {
131    pub fn from_env() -> Self {
132        let api_key = first_env(&[
133            "VITE_OPENRTC_KEY",
134            "VITE_OPENRTC_API_KEY",
135            "VITE_PLUTO_OPENRTC_API_KEY",
136            "OPENRTC_API_KEY",
137        ]);
138        Self {
139            api_key: api_key.unwrap_or_default(),
140            data_dir: None,
141            transport_config: None,
142        }
143    }
144
145    pub fn validated_api_key(&self) -> Result<&str, String> {
146        openrtc::validate_api_key(&self.api_key)
147            .map_err(|error| format!("invalid OpenRTC 2.0 public API key: {error}"))
148    }
149
150    pub fn app_tag(&self) -> Result<String, String> {
151        self.validated_api_key().map(openrtc::app_tag_from_api_key)
152    }
153}
154
155fn first_env(names: &[&str]) -> Option<String> {
156    names.iter().find_map(|name| {
157        std::env::var(name)
158            .ok()
159            .map(|value| value.trim().to_string())
160            .filter(|value| !value.is_empty())
161    })
162}
163
164#[derive(Default)]
165struct TokenRelayState {
166    identity_credential: RwLock<Option<String>>,
167}
168
169impl TokenRelayState {
170    fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
171        let relay = self.clone();
172        Box::new(move || {
173            relay
174                .identity_credential
175                .read()
176                .ok()
177                .and_then(|guard| guard.clone())
178        })
179    }
180
181    fn set(&self, identity_credential: Option<String>) {
182        if let Ok(mut guard) = self.identity_credential.write() {
183            *guard = normalize_token(identity_credential);
184        }
185    }
186}
187
188fn normalize_token(value: Option<String>) -> Option<String> {
189    value
190        .map(|value| value.trim().to_string())
191        .filter(|value| !value.is_empty())
192}
193
194#[derive(Debug, Clone)]
195struct NativeCapabilityRegistration {
196    avenue_kind: String,
197    avenue_id: String,
198    desired_revision: u64,
199    desired_peers: Vec<serde_json::Value>,
200}
201
202#[derive(Debug, Default)]
203struct NativeCapabilityRegistry {
204    registrations: HashMap<String, NativeCapabilityRegistration>,
205    root_desired_revision: u64,
206}
207
208impl NativeCapabilityRegistry {
209    fn register(
210        &mut self,
211        capability_key: String,
212        avenue_kind: String,
213        avenue_id: String,
214    ) -> Result<(), String> {
215        if let Some(existing) = self.registrations.get(&capability_key) {
216            if existing.avenue_kind == avenue_kind && existing.avenue_id == avenue_id {
217                return Ok(());
218            }
219            return Err(format!(
220                "native capability key {capability_key} is already registered for another avenue"
221            ));
222        }
223        self.registrations.insert(
224            capability_key,
225            NativeCapabilityRegistration {
226                avenue_kind,
227                avenue_id,
228                desired_revision: 0,
229                desired_peers: Vec::new(),
230            },
231        );
232        Ok(())
233    }
234
235    fn capability_keys_for_identity(
236        &self,
237        connection_id: Option<&str>,
238        device_id: Option<&str>,
239        device_id_hint: Option<&str>,
240        remote_node_id: Option<&str>,
241    ) -> Vec<String> {
242        let identities = [connection_id, device_id, device_id_hint, remote_node_id]
243            .into_iter()
244            .flatten()
245            .map(str::trim)
246            .filter(|value| !value.is_empty())
247            .collect::<BTreeSet<_>>();
248        self.registrations
249            .iter()
250            .filter_map(|(key, registration)| {
251                registration
252                    .desired_peers
253                    .iter()
254                    .any(|peer| {
255                        ["connectionId", "deviceId", "nodeId"]
256                            .into_iter()
257                            .filter_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
258                            .map(str::trim)
259                            .any(|value| identities.contains(value))
260                    })
261                    .then(|| key.clone())
262            })
263            .collect::<BTreeSet<_>>()
264            .into_iter()
265            .collect()
266    }
267
268    fn capability_keys_for_state(&self, snapshot: &openrtc::client::StateSnapshot) -> Vec<String> {
269        self.capability_keys_for_identity(
270            Some(&snapshot.connection_id),
271            snapshot.device_id.as_deref(),
272            snapshot.device_id_hint.as_deref(),
273            snapshot.remote_node_id.as_deref(),
274        )
275    }
276
277    fn capability_keys_for_peer_data(
278        &self,
279        event: &openrtc::client::NativePeerDataEvent,
280    ) -> Vec<String> {
281        if let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) {
282            if let Some(capability) = value
283                .get("capability")
284                .and_then(serde_json::Value::as_str)
285                .map(str::trim)
286                .filter(|value| !value.is_empty())
287            {
288                if self.registrations.contains_key(capability) {
289                    return vec![capability.to_string()];
290                }
291                return Vec::new();
292            }
293        }
294        self.capability_keys_for_identity(
295            Some(&event.connection_id),
296            None,
297            None,
298            event.remote_node_id.as_deref(),
299        )
300    }
301
302    fn projected_stream_capability(
303        &self,
304        channel: &openrtc::stream_metadata::ChannelMetadata,
305    ) -> Option<String> {
306        let explicit = channel
307            .metadata
308            .as_ref()
309            .and_then(|metadata| metadata.get("openrtcCapability"))
310            .and_then(serde_json::Value::as_str)
311            .map(str::trim)
312            .filter(|value| !value.is_empty());
313        match explicit {
314            Some(key) if self.registrations.contains_key(key) => Some(key.to_string()),
315            _ => None,
316        }
317    }
318
319    fn aggregate_desired_peers(&mut self) -> Result<(u64, String), String> {
320        let mut peers = std::collections::BTreeMap::<String, serde_json::Value>::new();
321        for registration in self.registrations.values() {
322            for peer in &registration.desired_peers {
323                let identity = ["deviceId", "nodeId", "ticket"]
324                    .into_iter()
325                    .find_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
326                    .map(str::trim)
327                    .filter(|value| !value.is_empty())
328                    .ok_or_else(|| "desired peer has no stable identity".to_string())?;
329                peers
330                    .entry(identity.to_string())
331                    .or_insert_with(|| peer.clone());
332            }
333        }
334        self.root_desired_revision = self.root_desired_revision.saturating_add(1);
335        let payload = serde_json::to_string(&peers.into_values().collect::<Vec<_>>())
336            .map_err(|error| format!("serialize aggregated desired peers: {error}"))?;
337        Ok((self.root_desired_revision, payload))
338    }
339}
340
341fn required_native_capability_part(value: String, label: &str) -> Result<String, String> {
342    let value = value.trim().to_string();
343    if value.is_empty()
344        || value.len() > 192
345        || !value
346            .chars()
347            .all(|character| character.is_ascii_alphanumeric() || "_.:@-".contains(character))
348    {
349        return Err(format!("native {label} is invalid"));
350    }
351    Ok(value)
352}
353
354pub struct OpenRtcTauriState {
355    client: RwLock<Arc<openrtc::client::Client>>,
356    config: OpenRtcTauriConfig,
357    token_relay: Arc<TokenRelayState>,
358    data_dir: Option<PathBuf>,
359    connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
360    peer_data_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
361    presence_loop_active: Mutex<bool>,
362    subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
363    native_projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
364    projected_stream_offered: AtomicU64,
365    projected_stream_decoded: AtomicU64,
366    projected_stream_projected: AtomicU64,
367    projected_stream_unhandled_non_channel: AtomicU64,
368    projected_stream_unhandled_other_channel: AtomicU64,
369    projected_stream_unhandled_no_subscription: AtomicU64,
370    projected_stream_failures: AtomicU64,
371    projected_stream_last_transport_stable_id: AtomicU64,
372    projected_stream_validation_checks: AtomicU64,
373    projected_stream_validation_rejections: AtomicU64,
374    projected_stream_last_validated_transport_stable_id: AtomicU64,
375    projected_stream_authorized: AtomicU64,
376    projected_stream_unauthorized: AtomicU64,
377    projected_stream_last_channel: Mutex<Option<String>>,
378    projected_stream_last_protocol: Mutex<Option<String>>,
379    projected_stream_last_connection_id: Mutex<Option<String>>,
380    projected_stream_last_remote_node_id: Mutex<Option<String>>,
381    peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
382    peer_uni_streams: Mutex<HashMap<String, Arc<Mutex<Option<PeerSendStream>>>>>,
383    portable_media: Mutex<openrtc::media::PortableMediaSession>,
384    managed_session_start_guard: Mutex<()>,
385    managed_session_next_owner_epoch: AtomicU64,
386    managed_session: Mutex<Option<ManagedSessionRecord>>,
387    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
388    native_transport_installers: Vec<Arc<dyn TransportInstaller>>,
389    installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
390    native_device_key_signer: Option<Arc<dyn DeviceKeySigner>>,
391}
392
393impl OpenRtcTauriState {
394    pub fn new(config: OpenRtcTauriConfig) -> Self {
395        openrtc::ensure_rustls();
396
397        let token_relay = Arc::new(TokenRelayState::default());
398        let client = build_client(&config, &token_relay)
399            .expect("OpenRTC Tauri 2.0 requires a valid public API key");
400        let data_dir = config.data_dir.clone();
401
402        Self {
403            client: RwLock::new(client),
404            config,
405            token_relay,
406            data_dir,
407            connection_state_forwarder: std::sync::Mutex::new(None),
408            peer_data_forwarder: std::sync::Mutex::new(None),
409            presence_loop_active: Mutex::new(false),
410            subscriptions: Mutex::new(HashMap::new()),
411            native_projection_subscription: Arc::new(Mutex::new(None)),
412            projected_stream_offered: AtomicU64::new(0),
413            projected_stream_decoded: AtomicU64::new(0),
414            projected_stream_projected: AtomicU64::new(0),
415            projected_stream_unhandled_non_channel: AtomicU64::new(0),
416            projected_stream_unhandled_other_channel: AtomicU64::new(0),
417            projected_stream_unhandled_no_subscription: AtomicU64::new(0),
418            projected_stream_failures: AtomicU64::new(0),
419            projected_stream_last_transport_stable_id: AtomicU64::new(0),
420            projected_stream_validation_checks: AtomicU64::new(0),
421            projected_stream_validation_rejections: AtomicU64::new(0),
422            projected_stream_last_validated_transport_stable_id: AtomicU64::new(0),
423            projected_stream_authorized: AtomicU64::new(0),
424            projected_stream_unauthorized: AtomicU64::new(0),
425            projected_stream_last_channel: Mutex::new(None),
426            projected_stream_last_protocol: Mutex::new(None),
427            projected_stream_last_connection_id: Mutex::new(None),
428            projected_stream_last_remote_node_id: Mutex::new(None),
429            peer_bi_streams: Mutex::new(HashMap::new()),
430            peer_uni_streams: Mutex::new(HashMap::new()),
431            portable_media: Mutex::new(openrtc::media::PortableMediaSession::default()),
432            managed_session_start_guard: Mutex::new(()),
433            managed_session_next_owner_epoch: AtomicU64::new(0),
434            managed_session: Mutex::new(None),
435            capability_registry: Arc::new(Mutex::new(NativeCapabilityRegistry::default())),
436            native_transport_installers: Vec::new(),
437            installed_native_transports: Mutex::new(HashMap::new()),
438            native_device_key_signer: None,
439        }
440    }
441
442    pub fn with_native_transport_installer(
443        mut self,
444        installer: Arc<dyn TransportInstaller>,
445    ) -> Self {
446        self.native_transport_installers.push(installer);
447        self
448    }
449
450    pub fn with_native_device_key_signer(mut self, signer: Arc<dyn DeviceKeySigner>) -> Self {
451        self.native_device_key_signer = Some(signer);
452        self
453    }
454
455    fn device_public_key(&self, app_tag: &str) -> Result<serde_json::Value, String> {
456        let app_tag = required_app_tag(app_tag)?;
457        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
458            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
459        })?;
460        let value = signer.public_jwk(app_tag)?;
461        validate_public_device_jwk(&value)?;
462        Ok(value)
463    }
464
465    fn sign_device_proof(&self, app_tag: &str, challenge: &str) -> Result<String, String> {
466        let app_tag = required_app_tag(app_tag)?;
467        if challenge.is_empty() || challenge.len() > 2_048 {
468            return Err("OpenRTC device challenge is invalid".to_string());
469        }
470        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
471            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
472        })?;
473        let signature = signer.sign(app_tag, challenge.as_bytes())?;
474        if signature.len() != 64 {
475            return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
476        }
477        Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
478    }
479
480    fn delete_device_key(&self, app_tag: &str) -> Result<(), String> {
481        let app_tag = required_app_tag(app_tag)?;
482        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
483            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
484        })?;
485        signer.delete(app_tag)
486    }
487
488    fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
489        let app_tag = required_app_tag(app_tag)?;
490        let key = required_secure_record_key(key)?;
491        self.native_device_key_signer
492            .as_ref()
493            .ok_or_else(|| {
494                "OpenRTC 2.0 native certificate persistence requires a host secure store"
495                    .to_string()
496            })?
497            .read_secure_record(app_tag, key)
498    }
499
500    fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
501        let app_tag = required_app_tag(app_tag)?;
502        let key = required_secure_record_key(key)?;
503        if value.is_empty() || value.len() > 16 * 1024 {
504            return Err("OpenRTC secure record is invalid".to_string());
505        }
506        self.native_device_key_signer
507            .as_ref()
508            .ok_or_else(|| {
509                "OpenRTC 2.0 native certificate persistence requires a host secure store"
510                    .to_string()
511            })?
512            .write_secure_record(app_tag, key, value)
513    }
514
515    fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
516        let app_tag = required_app_tag(app_tag)?;
517        let key = required_secure_record_key(key)?;
518        self.native_device_key_signer
519            .as_ref()
520            .ok_or_else(|| {
521                "OpenRTC 2.0 native certificate persistence requires a host secure store"
522                    .to_string()
523            })?
524            .delete_secure_record(app_tag, key)
525    }
526
527    pub fn client(&self) -> Arc<openrtc::client::Client> {
528        self.client
529            .read()
530            .map(|guard| guard.clone())
531            .expect("OpenRTC Tauri client state is poisoned")
532    }
533
534    pub fn set_identity_credential(&self, identity_credential: Option<String>) {
535        self.token_relay.set(identity_credential);
536    }
537
538    fn allocate_managed_session_owner_epoch(&self) -> u64 {
539        self.managed_session_next_owner_epoch
540            .fetch_add(1, Ordering::SeqCst)
541            + 1
542    }
543
544    fn replace_connection_state_forwarder(&self, client: Arc<openrtc::client::Client>) {
545        let capability_registry = self.capability_registry.clone();
546        let projection_subscription = self.native_projection_subscription.clone();
547        let next = tauri::async_runtime::spawn(async move {
548            forward_connection_state_events(client, capability_registry, projection_subscription)
549                .await;
550        });
551        if let Ok(mut guard) = self.connection_state_forwarder.lock() {
552            if let Some(previous) = guard.replace(next) {
553                previous.abort();
554            }
555        }
556    }
557
558    fn replace_peer_data_forwarder(&self, client: Arc<openrtc::client::Client>) {
559        let capability_registry = self.capability_registry.clone();
560        let projection_subscription = self.native_projection_subscription.clone();
561        let next = tauri::async_runtime::spawn(async move {
562            forward_peer_data_events(client, capability_registry, projection_subscription).await;
563        });
564        if let Ok(mut guard) = self.peer_data_forwarder.lock() {
565            if let Some(previous) = guard.replace(next) {
566                previous.abort();
567            }
568        }
569    }
570}
571
572fn build_client(
573    config: &OpenRtcTauriConfig,
574    token_relay: &Arc<TokenRelayState>,
575) -> Result<Arc<openrtc::client::Client>, String> {
576    let mut builder = openrtc::client::Client::builder(
577        config.validated_api_key()?.to_string(),
578        token_relay.token_provider(),
579    )
580    .map_err(|error| error.to_string())?;
581    if let Some(transport_config) = config.transport_config.clone() {
582        builder = builder.transport_config(transport_config);
583    }
584    Ok(Arc::new(builder.build()))
585}
586
587fn requested_local_device_id(value: Option<&str>) -> Option<String> {
588    value
589        .map(str::trim)
590        .filter(|value| !value.is_empty())
591        .map(ToOwned::to_owned)
592}
593
594fn managed_session_device_id(
595    _requested: Option<&str>,
596    native_identity: &openrtc::native_device::NativeDeviceIdentity,
597) -> String {
598    // The persisted native identity is the only durable-device authority.
599    // Frontend IDs are aliases/display hints and must never create a second
600    // gateway device projection or presence lease for the same native node.
601    native_identity.device_id.clone()
602}
603
604#[cfg(test)]
605fn metadata_with_authoritative_device_id(
606    metadata: Option<String>,
607    device_id: &str,
608) -> Option<String> {
609    let device_id = device_id.trim();
610    if device_id.is_empty() {
611        return metadata;
612    }
613
614    let Some(raw_metadata) = metadata else {
615        return Some(serde_json::json!({ "deviceId": device_id }).to_string());
616    };
617
618    match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
619        Ok(serde_json::Value::Object(mut map)) => {
620            map.insert(
621                "deviceId".to_string(),
622                serde_json::Value::String(device_id.to_string()),
623            );
624            Some(serde_json::Value::Object(map).to_string())
625        }
626        _ => Some(
627            serde_json::json!({
628                "deviceId": device_id,
629                "metadata": raw_metadata,
630            })
631            .to_string(),
632        ),
633    }
634}
635
636struct PeerBiStreamHandle {
637    send: Arc<Mutex<Option<PeerSendStream>>>,
638    recv: Option<PeerRecvStream>,
639    read_task: Option<tokio::task::JoinHandle<()>>,
640}
641
642struct NativeProjectionSubscription {
643    request_id: String,
644    channel: Channel<NativeProjectionEvent>,
645}
646
647#[derive(Serialize, Clone)]
648#[serde(rename_all = "camelCase")]
649struct NativeProjectionEvent {
650    request_id: String,
651    kind: &'static str,
652    capability_keys: Vec<String>,
653    payload: serde_json::Value,
654}
655
656#[derive(Serialize)]
657#[serde(rename_all = "camelCase")]
658pub struct OpenBiResult {
659    stream_id: String,
660    connection_id: Option<String>,
661    remote_node_id: String,
662}
663
664#[derive(Serialize)]
665#[serde(rename_all = "camelCase")]
666pub struct OpenUniResult {
667    stream_id: String,
668    connection_id: Option<String>,
669    remote_node_id: String,
670}
671
672#[derive(Serialize, Clone)]
673#[serde(rename_all = "camelCase")]
674struct IncomingPeerBiStreamEvent {
675    request_id: String,
676    stream_id: String,
677    connection_id: Option<String>,
678    remote_node_id: String,
679    transport_stable_id: u64,
680    channel: Option<openrtc::stream_metadata::ChannelMetadata>,
681    application_authorized: bool,
682    capability_key: String,
683}
684
685#[derive(Debug, Clone, Serialize)]
686#[serde(rename_all = "camelCase")]
687struct ProjectedStreamDiagnostics {
688    subscription_active: bool,
689    offered: u64,
690    decoded: u64,
691    projected: u64,
692    unhandled_non_channel: u64,
693    unhandled_other_channel: u64,
694    unhandled_no_subscription: u64,
695    failures: u64,
696    last_transport_stable_id: u64,
697    validation_checks: u64,
698    validation_rejections: u64,
699    last_validated_transport_stable_id: u64,
700    authorized: u64,
701    unauthorized: u64,
702    last_channel: Option<String>,
703    last_protocol: Option<String>,
704    last_connection_id: Option<String>,
705    last_remote_node_id: Option<String>,
706}
707
708/// Result of offering one already-admitted, already-protected peer stream to
709/// the OpenRTC-owned native protocol projection. A host must continue its own
710/// protocol dispatch only for `Unhandled`; there is still exactly one owner of
711/// the underlying incoming-stream queue.
712pub enum ProjectedPeerBiHandoff {
713    Handled,
714    Unhandled {
715        send: PeerSendStream,
716        recv: PeerRecvStream,
717    },
718}
719
720#[derive(Serialize, Clone)]
721#[serde(rename_all = "camelCase")]
722struct PeerBiStreamClosedEvent {
723    r#type: &'static str,
724    error: Option<String>,
725}
726
727#[derive(Debug, Clone, Serialize)]
728#[serde(rename_all = "camelCase")]
729pub struct StartSessionResult {
730    local_node_id: String,
731    ticket_scope: Option<String>,
732    ticket: Option<String>,
733    presence_started: bool,
734    auto_connect_started: bool,
735    local_device: openrtc::native_device::NativeDeviceIdentity,
736}
737
738#[derive(Clone)]
739struct ManagedSessionRecord {
740    key: String,
741    owner_client: Arc<openrtc::client::Client>,
742    owner_epoch: u64,
743    result: StartSessionResult,
744}
745
746#[derive(Debug, Clone, Copy, PartialEq, Eq)]
747enum SessionDisposition {
748    Start,
749    Reuse,
750    Refresh,
751    Replace,
752}
753
754fn managed_session_disposition(
755    active_key: Option<&str>,
756    active_client_matches: bool,
757    active_presence_started: bool,
758    active_auto_connect_started: bool,
759    requested_key: &str,
760    requested_presence: bool,
761    requested_auto_connect: bool,
762) -> SessionDisposition {
763    let Some(active_key) = active_key else {
764        return SessionDisposition::Start;
765    };
766    if !active_client_matches || active_key != requested_key {
767        return SessionDisposition::Replace;
768    }
769    if (!requested_presence || active_presence_started)
770        && (!requested_auto_connect || active_auto_connect_started)
771    {
772        SessionDisposition::Reuse
773    } else {
774        SessionDisposition::Refresh
775    }
776}
777
778fn revokes_managed_session(scope: &str) -> bool {
779    scope.trim() == "user-device"
780}
781
782pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
783    init_with_state(OpenRtcTauriState::new(config))
784}
785
786pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
787    tauri::plugin::Builder::new(PLUGIN_NAME)
788        .setup(move |app, _api| {
789            let client = state.client();
790            let transport_config = state.config.transport_config.clone();
791            // Optional native transports must decorate the endpoint before the
792            // host app can initialize the shared Iroh client during its setup.
793            // Managed-session startup is too late for custom transports because
794            // Iroh endpoint hooks are immutable after bind/adoption.
795            let install_context = InstallContext {
796                data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
797            };
798            tauri::async_runtime::block_on(ensure_requested_native_transports(
799                &state,
800                &client,
801                transport_config.as_ref(),
802                &install_context,
803            ))
804            .map_err(std::io::Error::other)?;
805            state.replace_connection_state_forwarder(client.clone());
806            state.replace_peer_data_forwarder(client);
807            app.manage(state);
808            Ok(())
809        })
810        .invoke_handler(tauri::generate_handler![
811            openrtc_set_identity_credential,
812            openrtc_device_public_key,
813            openrtc_sign_device_proof,
814            openrtc_delete_device_key,
815            openrtc_read_secure_record,
816            openrtc_write_secure_record,
817            openrtc_delete_secure_record,
818            rtc_native_status,
819            get_rtc_local_device_info,
820            update_rtc_local_device_name,
821            get_iroh_node_id,
822            start_iroh_node,
823            get_iroh_endpoint_ticket,
824            register_session_token,
825            get_endpoint_ticket_with_token,
826            validate_session_token,
827            revoke_session_tokens_by_scope,
828            stop_rtc_presence_loop,
829            register_rtc_capability,
830            unregister_rtc_capability,
831            start_rtc_managed_session,
832            start_rtc_external_auto_connect,
833            submit_rtc_desired_peers,
834            stop_rtc_auto_connect,
835            notify_rtc_network_change,
836            set_rtc_transport_priority,
837            connect_to_device,
838            disconnect_device,
839            set_auto_connect_excluded,
840            set_rtc_external_auto_connect_excluded,
841            resolve_rtc_peer_connection_records,
842            resolve_rtc_peer_identity,
843            get_rtc_peer_session,
844            list_rtc_peer_sessions,
845            list_rtc_managed_connections,
846            wait_for_rtc_settled_peer,
847            list_rtc_connection_states,
848            get_rtc_connection_state,
849            stop_rtc_subscription,
850            start_native_projection,
851            get_projected_stream_diagnostics,
852            is_current_transport_stable_id,
853            send_peer_message,
854            is_peer_connected,
855            open_peer_bi_stream,
856            open_peer_bi_transport_only_stream,
857            open_peer_uni_stream,
858            write_peer_bi_stream,
859            start_peer_bi_stream_read,
860            cancel_peer_bi_stream_read,
861            close_peer_bi_stream,
862            write_peer_uni_stream,
863            close_peer_uni_stream,
864            begin_openrtc_media_publication,
865            encode_openrtc_media_sample,
866            pause_openrtc_media_publication,
867            retire_openrtc_media_publication,
868            retire_openrtc_media_receiver,
869            decode_openrtc_media_chunk,
870            decode_openrtc_media_control,
871        ])
872        .build()
873}
874
875pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
876    init(OpenRtcTauriConfig::from_env())
877}
878
879async fn send_native_projection<T: Serialize>(
880    subscription: &Arc<Mutex<Option<NativeProjectionSubscription>>>,
881    kind: &'static str,
882    capability_keys: Vec<String>,
883    payload: &T,
884) {
885    if capability_keys.is_empty() {
886        return;
887    }
888    let Ok(payload) = serde_json::to_value(payload) else {
889        return;
890    };
891    let target = subscription.lock().await.as_ref().map(|subscription| {
892        (
893            subscription.request_id.clone(),
894            subscription.channel.clone(),
895        )
896    });
897    let Some((request_id, channel)) = target else {
898        return;
899    };
900    let _ = channel.send(NativeProjectionEvent {
901        request_id,
902        kind,
903        capability_keys,
904        payload,
905    });
906}
907
908async fn forward_connection_state_events(
909    client: Arc<openrtc::client::Client>,
910    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
911    projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
912) {
913    let mut rx = client.connection_state_updates();
914    while let Ok(snapshot) = rx.recv().await {
915        let keys = capability_registry
916            .lock()
917            .await
918            .capability_keys_for_state(&snapshot);
919        send_native_projection(
920            &projection_subscription,
921            "connection-state",
922            keys,
923            &snapshot,
924        )
925        .await;
926    }
927}
928
929async fn forward_peer_data_events(
930    client: Arc<openrtc::client::Client>,
931    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
932    projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
933) {
934    let mut rx = client.subscribe_native_peer_data();
935    while let Ok(event) = rx.recv().await {
936        let keys = capability_registry
937            .lock()
938            .await
939            .capability_keys_for_peer_data(&event);
940        send_native_projection(&projection_subscription, "peer-data", keys, &event).await;
941    }
942}
943
944async fn replay_current_connection_states(state: &OpenRtcTauriState) {
945    let snapshots = state.client().connection_states().await;
946    let registry = state.capability_registry.lock().await;
947    for snapshot in snapshots {
948        let keys = registry.capability_keys_for_state(&snapshot);
949        send_native_projection(
950            &state.native_projection_subscription,
951            "connection-state",
952            keys,
953            &snapshot,
954        )
955        .await;
956    }
957}
958
959fn app_data_dir<R: Runtime>(
960    app: &tauri::AppHandle<R>,
961    state: &OpenRtcTauriState,
962) -> Result<PathBuf, String> {
963    state
964        .data_dir
965        .clone()
966        .or_else(|| app.path().app_data_dir().ok())
967        .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
968}
969
970async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
971    if let Some(node_id) = client.current_node_id().await {
972        return Ok(node_id);
973    }
974    client
975        .init_iroh(None, Vec::new())
976        .await
977        .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
978}
979
980async fn ensure_requested_native_transports(
981    state: &OpenRtcTauriState,
982    client: &Arc<openrtc::client::Client>,
983    transports: Option<&openrtc::client::TransportConfig>,
984    context: &InstallContext,
985) -> Result<(), String> {
986    let Some(transports) = transports else {
987        return Ok(());
988    };
989    let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
990    if ble_requested
991        && !state
992            .native_transport_installers
993            .iter()
994            .any(|installer| installer.is_requested(transports))
995    {
996        eprintln!(
997            "[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
998        );
999        return Ok(());
1000    }
1001    let client_ptr = Arc::as_ptr(client) as usize;
1002    for installer in &state.native_transport_installers {
1003        if !installer.is_requested(transports) {
1004            continue;
1005        }
1006        let already_installed = state
1007            .installed_native_transports
1008            .lock()
1009            .await
1010            .get(installer.id())
1011            .is_some_and(|installed| installed.client_ptr == client_ptr);
1012        if already_installed {
1013            continue;
1014        }
1015        if client.current_node_id().await.is_some() {
1016            eprintln!(
1017                "[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
1018                installer.id()
1019            );
1020            continue;
1021        }
1022        let runtime = match installer
1023            .install(client.clone(), transports.clone(), context.clone())
1024            .await
1025        {
1026            Ok(runtime) => runtime,
1027            Err(error) => {
1028                eprintln!(
1029                    "[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
1030                    installer.id(),
1031                    error
1032                );
1033                continue;
1034            }
1035        };
1036        state.installed_native_transports.lock().await.insert(
1037            installer.id(),
1038            InstalledNativeTransport {
1039                client_ptr,
1040                _runtime: runtime,
1041            },
1042        );
1043    }
1044    Ok(())
1045}
1046
1047async fn register_peer_bi_stream(
1048    state: &OpenRtcTauriState,
1049    connection_id: Option<String>,
1050    remote_node_id: String,
1051    send: PeerSendStream,
1052    recv: PeerRecvStream,
1053) -> OpenBiResult {
1054    let stream_id = uuid::Uuid::new_v4().to_string();
1055    let send = Arc::new(Mutex::new(Some(send)));
1056
1057    state.peer_bi_streams.lock().await.insert(
1058        stream_id.clone(),
1059        PeerBiStreamHandle {
1060            send,
1061            // Streams stay paused until JavaScript supplies its ordered Tauri
1062            // Channel. This prevents byte zero from racing listener setup.
1063            recv: Some(recv),
1064            read_task: None,
1065        },
1066    );
1067
1068    OpenBiResult {
1069        stream_id,
1070        connection_id,
1071        remote_node_id,
1072    }
1073}
1074
1075async fn read_peer_exact(recv: &mut PeerRecvStream, buffer: &mut [u8]) -> std::io::Result<()> {
1076    let mut offset = 0;
1077    while offset < buffer.len() {
1078        let read = recv.read(&mut buffer[offset..]).await?;
1079        if read == 0 {
1080            return Err(std::io::Error::new(
1081                std::io::ErrorKind::UnexpectedEof,
1082                "peer stream ended during channel classification",
1083            ));
1084        }
1085        offset += read;
1086    }
1087    Ok(())
1088}
1089
1090fn is_openrtc_projected_channel(channel: &openrtc::stream_metadata::ChannelMetadata) -> bool {
1091    channel.channel_id == openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID
1092}
1093
1094/// Offer an already-admitted and application-protected stream from the
1095/// embedding native host to OpenRTC's typed IPC projection.
1096///
1097/// This is the integration point for Tauri applications that already consume
1098/// the client's sole `incoming_streams()` queue for product protocols. The
1099/// helper performs bounded channel-envelope classification and restores every
1100/// consumed plaintext byte when the stream is not owned by OpenRTC.
1101pub async fn handoff_projected_peer_bi_stream<R: Runtime>(
1102    app: tauri::AppHandle<R>,
1103    connection_id: Option<String>,
1104    remote_node_id: String,
1105    transport_stable_id: u64,
1106    send: PeerSendStream,
1107    mut recv: PeerRecvStream,
1108) -> Result<ProjectedPeerBiHandoff, String> {
1109    let state = app.state::<OpenRtcTauriState>();
1110    state
1111        .projected_stream_offered
1112        .fetch_add(1, Ordering::Relaxed);
1113    state
1114        .projected_stream_last_transport_stable_id
1115        .store(transport_stable_id, Ordering::Relaxed);
1116    *state.projected_stream_last_connection_id.lock().await = connection_id.clone();
1117    *state.projected_stream_last_remote_node_id.lock().await = Some(remote_node_id.clone());
1118    let mut envelope = Vec::new();
1119    let channel = loop {
1120        match openrtc::stream_metadata::decode_prefix(&envelope) {
1121            openrtc::stream_metadata::DecodeDecision::NotMatched => {
1122                state
1123                    .projected_stream_unhandled_non_channel
1124                    .fetch_add(1, Ordering::Relaxed);
1125                return Ok(ProjectedPeerBiHandoff::Unhandled {
1126                    send,
1127                    recv: recv.with_plaintext_prefix(envelope),
1128                });
1129            }
1130            openrtc::stream_metadata::DecodeDecision::Decoded(prefix) => break prefix.channel,
1131            openrtc::stream_metadata::DecodeDecision::NeedMore(required) => {
1132                let missing = required.saturating_sub(envelope.len());
1133                if missing == 0 {
1134                    return Err("channel-envelope decoder made no progress".to_string());
1135                }
1136                let mut bytes = vec![0u8; missing];
1137                if let Err(error) = read_peer_exact(&mut recv, &mut bytes).await {
1138                    state
1139                        .projected_stream_failures
1140                        .fetch_add(1, Ordering::Relaxed);
1141                    return Err(format!("read projected channel envelope: {error}"));
1142                }
1143                envelope.extend_from_slice(&bytes);
1144            }
1145        }
1146    };
1147    state
1148        .projected_stream_decoded
1149        .fetch_add(1, Ordering::Relaxed);
1150    *state.projected_stream_last_channel.lock().await = Some(channel.channel_id.clone());
1151    *state.projected_stream_last_protocol.lock().await = channel
1152        .metadata
1153        .as_ref()
1154        .and_then(|metadata| metadata.get("protocol"))
1155        .and_then(serde_json::Value::as_str)
1156        .map(str::to_string);
1157
1158    if !is_openrtc_projected_channel(&channel) {
1159        state
1160            .projected_stream_unhandled_other_channel
1161            .fetch_add(1, Ordering::Relaxed);
1162        return Ok(ProjectedPeerBiHandoff::Unhandled {
1163            send,
1164            recv: recv.with_plaintext_prefix(envelope),
1165        });
1166    }
1167
1168    let Some(capability_key) = state
1169        .capability_registry
1170        .lock()
1171        .await
1172        .projected_stream_capability(&channel)
1173    else {
1174        // Capability ownership is explicit on the protected channel envelope.
1175        // Unknown or unscoped descriptors remain available to the embedding
1176        // host; they are never assigned by listener order or capability count.
1177        state
1178            .projected_stream_unhandled_other_channel
1179            .fetch_add(1, Ordering::Relaxed);
1180        return Ok(ProjectedPeerBiHandoff::Unhandled {
1181            send,
1182            recv: recv.with_plaintext_prefix(envelope),
1183        });
1184    };
1185
1186    let subscription = state.native_projection_subscription.lock().await;
1187    let Some((request_id, projection_channel)) = subscription.as_ref().map(|subscription| {
1188        (
1189            subscription.request_id.clone(),
1190            subscription.channel.clone(),
1191        )
1192    }) else {
1193        state
1194            .projected_stream_unhandled_no_subscription
1195            .fetch_add(1, Ordering::Relaxed);
1196        return Ok(ProjectedPeerBiHandoff::Unhandled {
1197            send,
1198            recv: recv.with_plaintext_prefix(envelope),
1199        });
1200    };
1201    drop(subscription);
1202
1203    let application_authorized = connection_id.as_deref().is_some_and(|connection_id| {
1204        matches!(
1205            state.client().session_admission(connection_id),
1206            openrtc::session_token::SessionAdmission::Accepted { .. }
1207        )
1208    });
1209    if application_authorized {
1210        state
1211            .projected_stream_authorized
1212            .fetch_add(1, Ordering::Relaxed);
1213    } else {
1214        state
1215            .projected_stream_unauthorized
1216            .fetch_add(1, Ordering::Relaxed);
1217    }
1218
1219    let result = register_peer_bi_stream(
1220        state.inner(),
1221        connection_id,
1222        remote_node_id.clone(),
1223        send,
1224        recv,
1225    )
1226    .await;
1227    let event = IncomingPeerBiStreamEvent {
1228        request_id: request_id.clone(),
1229        stream_id: result.stream_id.clone(),
1230        connection_id: result.connection_id,
1231        remote_node_id,
1232        transport_stable_id,
1233        channel: Some(channel),
1234        application_authorized,
1235        capability_key: capability_key.clone(),
1236    };
1237    let payload = serde_json::to_value(event)
1238        .map_err(|error| format!("serialize projected peer stream: {error}"))?;
1239    if let Err(error) = projection_channel.send(NativeProjectionEvent {
1240        request_id,
1241        kind: "stream",
1242        capability_keys: vec![capability_key],
1243        payload,
1244    }) {
1245        state.peer_bi_streams.lock().await.remove(&result.stream_id);
1246        state
1247            .projected_stream_failures
1248            .fetch_add(1, Ordering::Relaxed);
1249        return Err(format!("send projected peer stream: {error}"));
1250    }
1251    state
1252        .projected_stream_projected
1253        .fetch_add(1, Ordering::Relaxed);
1254    Ok(ProjectedPeerBiHandoff::Handled)
1255}
1256
1257fn native_stream_trace_enabled() -> bool {
1258    cfg!(debug_assertions)
1259        || std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1")
1260}
1261
1262fn spawn_peer_bi_stream_reader(
1263    channel: Channel<Response>,
1264    stream_id: String,
1265    mut recv: PeerRecvStream,
1266) -> tokio::task::JoinHandle<()> {
1267    tokio::spawn(async move {
1268        let mut chunk = vec![0_u8; 64 * 1024];
1269        let mut close_error: Option<String> = None;
1270        loop {
1271            match recv.read(&mut chunk).await {
1272                Ok(0) => break,
1273                Ok(n) => {
1274                    if native_stream_trace_enabled() {
1275                        eprintln!(
1276                            "[openrtc-tauri][peer-bi-stream] stream_id={} phase=plaintext-chunk bytes={}",
1277                            stream_id, n
1278                        );
1279                    }
1280                    if channel.send(Response::new(chunk[..n].to_vec())).is_err() {
1281                        close_error = Some("native stream IPC channel closed".to_string());
1282                        break;
1283                    }
1284                }
1285                Err(error) => {
1286                    if native_stream_trace_enabled() {
1287                        eprintln!(
1288                            "[openrtc-tauri][peer-bi-stream] stream_id={} phase=read-error error={}",
1289                            stream_id, error
1290                        );
1291                    }
1292                    close_error = Some(error.to_string());
1293                    break;
1294                }
1295            }
1296        }
1297        let close = serde_json::to_string(&PeerBiStreamClosedEvent {
1298            r#type: "closed",
1299            error: close_error,
1300        })
1301        .expect("peer stream close event serializes");
1302        let _ = channel.send(Response::new(close));
1303    })
1304}
1305
1306async fn open_peer_bi_with<R: Runtime, F, Fut>(
1307    _app: tauri::AppHandle<R>,
1308    state: tauri::State<'_, OpenRtcTauriState>,
1309    peer_id: String,
1310    timeout_ms: Option<u64>,
1311    open: F,
1312) -> Result<OpenBiResult, String>
1313where
1314    F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
1315    Fut: std::future::Future<
1316        Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
1317    >,
1318{
1319    let peer_id = peer_id.trim().to_string();
1320    if peer_id.is_empty() {
1321        return Err("peerId is required".to_string());
1322    }
1323
1324    let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
1325        .await
1326        .map_err(|error| format!("open peer bi stream failed: {error}"))?;
1327    Ok(register_peer_bi_stream(&state, connection_id, remote_node_id, send, recv).await)
1328}
1329
1330/// Relay an opaque OpenRTC-issued identity credential into the provider-neutral
1331/// Rust runtime. The host owns acquisition and renewal; this command never
1332/// interprets the value as a Firebase token and accepts no refresh token.
1333#[tauri::command]
1334async fn openrtc_set_identity_credential(
1335    state: tauri::State<'_, OpenRtcTauriState>,
1336    identity_credential: Option<String>,
1337) -> Result<(), String> {
1338    state.set_identity_credential(identity_credential);
1339    if *state.presence_loop_active.lock().await {
1340        state.client().request_presence_update();
1341    }
1342    Ok(())
1343}
1344
1345fn required_app_tag(value: &str) -> Result<&str, String> {
1346    let value = value.trim();
1347    if value.len() < 5
1348        || value.len() > 80
1349        || !value.bytes().all(|byte| {
1350            byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
1351        })
1352    {
1353        return Err("OpenRTC app tag is invalid".to_string());
1354    }
1355    Ok(value)
1356}
1357
1358fn validate_public_device_jwk(value: &serde_json::Value) -> Result<(), String> {
1359    let object = value
1360        .as_object()
1361        .ok_or_else(|| "OpenRTC device public key is invalid".to_string())?;
1362    let x = object
1363        .get("x")
1364        .and_then(serde_json::Value::as_str)
1365        .unwrap_or_default();
1366    if object.get("kty").and_then(serde_json::Value::as_str) != Some("OKP")
1367        || object.get("crv").and_then(serde_json::Value::as_str) != Some("Ed25519")
1368        || object.contains_key("d")
1369        || x.len() != 43
1370        || !x
1371            .bytes()
1372            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
1373    {
1374        return Err("OpenRTC device public key must be a public Ed25519 JWK".to_string());
1375    }
1376    Ok(())
1377}
1378
1379fn required_secure_record_key(value: &str) -> Result<&str, String> {
1380    let value = value.trim();
1381    if value.is_empty()
1382        || value.len() > 512
1383        || !value
1384            .bytes()
1385            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':'))
1386    {
1387        return Err("OpenRTC secure record key is invalid".to_string());
1388    }
1389    Ok(value)
1390}
1391
1392#[tauri::command]
1393async fn openrtc_device_public_key(
1394    state: tauri::State<'_, OpenRtcTauriState>,
1395    app_tag: String,
1396) -> Result<serde_json::Value, String> {
1397    state.device_public_key(&app_tag)
1398}
1399
1400#[tauri::command]
1401async fn openrtc_sign_device_proof(
1402    state: tauri::State<'_, OpenRtcTauriState>,
1403    app_tag: String,
1404    challenge: String,
1405) -> Result<String, String> {
1406    state.sign_device_proof(&app_tag, &challenge)
1407}
1408
1409#[tauri::command]
1410async fn openrtc_delete_device_key(
1411    state: tauri::State<'_, OpenRtcTauriState>,
1412    app_tag: String,
1413) -> Result<(), String> {
1414    state.delete_device_key(&app_tag)
1415}
1416
1417#[tauri::command]
1418async fn openrtc_read_secure_record(
1419    state: tauri::State<'_, OpenRtcTauriState>,
1420    app_tag: String,
1421    key: String,
1422) -> Result<Option<String>, String> {
1423    state.read_secure_record(&app_tag, &key)
1424}
1425
1426#[tauri::command]
1427async fn openrtc_write_secure_record(
1428    state: tauri::State<'_, OpenRtcTauriState>,
1429    app_tag: String,
1430    key: String,
1431    value: String,
1432) -> Result<(), String> {
1433    state.write_secure_record(&app_tag, &key, &value)
1434}
1435
1436#[tauri::command]
1437async fn openrtc_delete_secure_record(
1438    state: tauri::State<'_, OpenRtcTauriState>,
1439    app_tag: String,
1440    key: String,
1441) -> Result<(), String> {
1442    state.delete_secure_record(&app_tag, &key)
1443}
1444
1445#[tauri::command]
1446async fn rtc_native_status(
1447    state: tauri::State<'_, OpenRtcTauriState>,
1448) -> Result<openrtc::client::RuntimeStatus, String> {
1449    Ok(state.client().runtime_status().await)
1450}
1451
1452#[tauri::command]
1453async fn get_rtc_local_device_info<R: Runtime>(
1454    app: tauri::AppHandle<R>,
1455    state: tauri::State<'_, OpenRtcTauriState>,
1456) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1457    let data_dir = app_data_dir(&app, &state)?;
1458    state
1459        .client()
1460        .init_native_device_identity(data_dir, None)
1461        .await
1462        .map_err(|error| error.to_string())
1463}
1464
1465#[tauri::command]
1466async fn update_rtc_local_device_name<R: Runtime>(
1467    app: tauri::AppHandle<R>,
1468    state: tauri::State<'_, OpenRtcTauriState>,
1469    device_name: String,
1470) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1471    let data_dir = app_data_dir(&app, &state)?;
1472    let _ = state
1473        .client()
1474        .init_native_device_identity(data_dir, None)
1475        .await
1476        .map_err(|error| error.to_string())?;
1477    state
1478        .client()
1479        .update_native_device_name(&device_name)
1480        .await
1481        .map_err(|error| error.to_string())
1482}
1483
1484#[tauri::command]
1485async fn get_iroh_node_id(
1486    state: tauri::State<'_, OpenRtcTauriState>,
1487) -> Result<Option<String>, String> {
1488    Ok(state.client().current_node_id().await)
1489}
1490
1491#[tauri::command]
1492async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
1493    ensure_iroh_node(state.client()).await
1494}
1495
1496#[tauri::command]
1497async fn get_iroh_endpoint_ticket(
1498    state: tauri::State<'_, OpenRtcTauriState>,
1499) -> Result<String, String> {
1500    ensure_iroh_node(state.client()).await?;
1501    state
1502        .client()
1503        .endpoint_ticket()
1504        .await
1505        .map_err(|error| error.to_string())
1506}
1507
1508#[tauri::command]
1509async fn register_session_token(
1510    state: tauri::State<'_, OpenRtcTauriState>,
1511    token: String,
1512    scope: String,
1513    max_connections: u32,
1514    expires_at_ms: Option<u64>,
1515) -> Result<(), String> {
1516    if let Some(expires_at_ms) = expires_at_ms {
1517        state
1518            .client()
1519            .register_token_until(token, scope, max_connections, expires_at_ms);
1520    } else {
1521        state
1522            .client()
1523            .register_session_token(token, scope, max_connections);
1524    }
1525    Ok(())
1526}
1527
1528#[tauri::command]
1529async fn get_endpoint_ticket_with_token(
1530    state: tauri::State<'_, OpenRtcTauriState>,
1531    scope: String,
1532    max_connections: u32,
1533) -> Result<String, String> {
1534    ensure_iroh_node(state.client()).await?;
1535    state
1536        .client()
1537        .endpoint_ticket_with_token(&scope, max_connections)
1538        .await
1539        .map_err(|error| error.to_string())
1540}
1541
1542#[tauri::command]
1543async fn validate_session_token(
1544    state: tauri::State<'_, OpenRtcTauriState>,
1545    token: String,
1546    connection_id: Option<String>,
1547) -> Result<String, String> {
1548    if let Some(connection_id) = connection_id
1549        .as_deref()
1550        .map(str::trim)
1551        .filter(|value| !value.is_empty())
1552    {
1553        state
1554            .client()
1555            .validate_connection_token(&token, connection_id, None)
1556            .await
1557            .map_err(|error| error.to_string())
1558    } else {
1559        state.client().validate_session_token(&token)
1560    }
1561}
1562
1563#[tauri::command]
1564async fn revoke_session_tokens_by_scope(
1565    state: tauri::State<'_, OpenRtcTauriState>,
1566    scope: String,
1567) -> Result<Vec<String>, String> {
1568    let affected = state.client().revoke_tokens_by_scope(&scope).await;
1569    if revokes_managed_session(&scope) {
1570        // The cached managed-session result contains the compound ticket whose
1571        // token was just revoked. Keeping that record would make the next
1572        // same-user/device start return a credential the host can no longer
1573        // validate. Invalidate only this lifecycle cache; the durable native
1574        // device identity remains intact so explicit reauthentication can
1575        // re-enroll the same installation with a freshly minted ticket.
1576        if let Some(previous) = state.managed_session.lock().await.take() {
1577            previous.owner_client.stop_presence_loop();
1578            previous.owner_client.stop_auto_connect();
1579            previous.owner_client.stop_external_auto_connect().await;
1580        }
1581        *state.presence_loop_active.lock().await = false;
1582    }
1583    Ok(affected)
1584}
1585
1586#[tauri::command]
1587async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1588    state.client().stop_presence_loop();
1589    *state.presence_loop_active.lock().await = false;
1590    if let Some(active) = state.managed_session.lock().await.as_mut() {
1591        active.result.presence_started = false;
1592    }
1593    stop_subscription_by_id(&state, "internal-presence-loop").await;
1594    Ok(())
1595}
1596
1597/// Register one public capability with the Rust-owned native root.
1598///
1599/// Owner: openrtc-tauri-plugin native root registry.
1600/// Consumers: OpenRTC native capability adapters.
1601/// Introduced: 2026-08-27.
1602#[tauri::command]
1603async fn register_rtc_capability(
1604    state: tauri::State<'_, OpenRtcTauriState>,
1605    capability_key: String,
1606    avenue_kind: String,
1607    avenue_id: String,
1608) -> Result<(), String> {
1609    let capability_key = required_native_capability_part(capability_key, "capability key")?;
1610    let avenue_kind = required_native_capability_part(avenue_kind, "avenue kind")?;
1611    let avenue_id = required_native_capability_part(avenue_id, "avenue id")?;
1612    if capability_key != format!("{avenue_kind}:{avenue_id}") {
1613        return Err("native capability key does not match its avenue".to_string());
1614    }
1615    state
1616        .capability_registry
1617        .lock()
1618        .await
1619        .register(capability_key, avenue_kind, avenue_id)
1620}
1621
1622#[tauri::command]
1623async fn unregister_rtc_capability(
1624    state: tauri::State<'_, OpenRtcTauriState>,
1625    capability_key: String,
1626) -> Result<(), String> {
1627    let capability_key = required_native_capability_part(capability_key, "capability key")?;
1628    let (empty, aggregate) = {
1629        let mut registry = state.capability_registry.lock().await;
1630        registry.registrations.remove(&capability_key);
1631        let empty = registry.registrations.is_empty();
1632        let aggregate = (!empty)
1633            .then(|| registry.aggregate_desired_peers())
1634            .transpose()?;
1635        (empty, aggregate)
1636    };
1637    if empty {
1638        state.client().stop_external_auto_connect().await;
1639    } else if let Some((revision, peers_json)) = aggregate {
1640        state
1641            .client()
1642            .submit_external_desired_peers(revision, &peers_json)
1643            .await
1644            .map_err(|error| error.to_string())?;
1645    }
1646    Ok(())
1647}
1648
1649#[tauri::command]
1650async fn start_rtc_managed_session<R: Runtime>(
1651    app: tauri::AppHandle<R>,
1652    state: tauri::State<'_, OpenRtcTauriState>,
1653    user_id: String,
1654    device_name: Option<String>,
1655    local_device_id: Option<String>,
1656    metadata: Option<String>,
1657    transports: Option<openrtc::client::TransportConfig>,
1658    auto_connect: Option<bool>,
1659    presence: Option<bool>,
1660) -> Result<StartSessionResult, String> {
1661    let command_started = std::time::Instant::now();
1662    let _start_guard = state.managed_session_start_guard.lock().await;
1663    let client = state.client();
1664    let app_tag = client.app_tag().to_string();
1665    eprintln!(
1666        "[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
1667        user_id,
1668        app_tag,
1669        device_name
1670            .as_deref()
1671            .map(str::trim)
1672            .map(|value| !value.is_empty())
1673            .unwrap_or(false),
1674        local_device_id
1675            .as_deref()
1676            .map(str::trim)
1677            .map(|value| !value.is_empty())
1678            .unwrap_or(false),
1679        metadata
1680            .as_deref()
1681            .map(str::trim)
1682            .map(|value| !value.is_empty())
1683            .unwrap_or(false),
1684        auto_connect.unwrap_or(false),
1685        presence.unwrap_or(false)
1686    );
1687    let data_dir = app_data_dir(&app, &state)?;
1688    let install_context = InstallContext {
1689        data_dir: data_dir.clone(),
1690    };
1691    ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
1692        .await?;
1693    if let Some(transport_config) = transports {
1694        client
1695            .update_transport_config(transport_config)
1696            .await
1697            .map_err(|error| error.to_string())?;
1698    }
1699
1700    let local_node_id = ensure_iroh_node(client.clone()).await?;
1701    let local_device = client
1702        .init_native_device_identity(data_dir, device_name.as_deref())
1703        .await
1704        .map_err(|error| error.to_string())?;
1705
1706    let effective_local_device_id =
1707        managed_session_device_id(local_device_id.as_deref(), &local_device);
1708    if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
1709        if requested != local_device.device_id {
1710            eprintln!(
1711                "[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
1712                requested, local_device.device_id
1713            );
1714        }
1715    }
1716
1717    let should_start_presence = presence.unwrap_or(false);
1718    let should_start_auto_connect = auto_connect.unwrap_or(false);
1719    if should_start_presence || should_start_auto_connect {
1720        return Err(
1721            "native Firebase coordination has been removed; publish gateway presence and submit external desired peers instead"
1722                .to_string(),
1723        );
1724    }
1725    // The native endpoint and durable device are root-owned. Capability
1726    // principals live in the registry above and must not replace this physical
1727    // lifecycle owner when devices, spaces, rooms, or tickets coexist.
1728    let managed_session_key = format!("{}:{}", app_tag, effective_local_device_id);
1729    let (disposition, reused) = {
1730        let active = state.managed_session.lock().await;
1731        let disposition = managed_session_disposition(
1732            active.as_ref().map(|record| record.key.as_str()),
1733            active
1734                .as_ref()
1735                .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1736            active
1737                .as_ref()
1738                .is_some_and(|record| record.result.presence_started),
1739            active
1740                .as_ref()
1741                .is_some_and(|record| record.result.auto_connect_started),
1742            &managed_session_key,
1743            should_start_presence,
1744            should_start_auto_connect,
1745        );
1746        let reused = (disposition == SessionDisposition::Reuse).then(|| {
1747            active
1748                .as_ref()
1749                .expect("managed session reuse requires an active record")
1750                .result
1751                .clone()
1752        });
1753        (disposition, reused)
1754    };
1755    if let Some(reused) = reused {
1756        eprintln!(
1757            "[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
1758            managed_session_key,
1759            command_started.elapsed().as_millis()
1760        );
1761        return Ok(reused);
1762    }
1763
1764    // A different app/user/device tuple replaces the previous lifecycle owner.
1765    // Stop old background work before publishing the new authoritative lease.
1766    if disposition == SessionDisposition::Replace {
1767        let previous = state
1768            .managed_session
1769            .lock()
1770            .await
1771            .take()
1772            .expect("managed session replacement requires an active record");
1773        previous.owner_client.stop_presence_loop();
1774        previous.owner_client.stop_auto_connect();
1775        previous.owner_client.stop_external_auto_connect().await;
1776        *state.presence_loop_active.lock().await = false;
1777    }
1778    let ticket_scope = Some("user-device".to_string());
1779    let ticket = client
1780        .endpoint_ticket_with_token("user-device", 0)
1781        .await
1782        .map_err(|error| error.to_string())?;
1783
1784    let logical_local_device = local_device.clone();
1785
1786    let owner_epoch = match disposition {
1787        SessionDisposition::Refresh => {
1788            state
1789                .managed_session
1790                .lock()
1791                .await
1792                .as_ref()
1793                .expect("managed session refresh requires an active record")
1794                .owner_epoch
1795        }
1796        SessionDisposition::Start | SessionDisposition::Replace => {
1797            state.allocate_managed_session_owner_epoch()
1798        }
1799        SessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
1800    };
1801
1802    eprintln!(
1803        "[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={}",
1804        user_id,
1805        app_tag,
1806        local_node_id,
1807        logical_local_device.device_id,
1808        owner_epoch,
1809        should_start_presence,
1810        should_start_auto_connect,
1811        ticket_scope.as_deref().unwrap_or("unrestricted"),
1812        command_started.elapsed().as_millis()
1813    );
1814
1815    let result = StartSessionResult {
1816        local_node_id,
1817        ticket_scope,
1818        ticket: Some(ticket),
1819        presence_started: false,
1820        auto_connect_started: false,
1821        local_device: logical_local_device,
1822    };
1823    *state.managed_session.lock().await = Some(ManagedSessionRecord {
1824        key: managed_session_key,
1825        owner_client: client,
1826        owner_epoch,
1827        result: result.clone(),
1828    });
1829    Ok(result)
1830}
1831
1832#[tauri::command]
1833async fn start_rtc_external_auto_connect(
1834    state: tauri::State<'_, OpenRtcTauriState>,
1835    user_id: String,
1836    local_device_id: String,
1837    capability_key: Option<String>,
1838) -> Result<(), String> {
1839    if let Some(capability_key) = capability_key {
1840        let capability_key = required_native_capability_part(capability_key, "capability key")?;
1841        if !state
1842            .capability_registry
1843            .lock()
1844            .await
1845            .registrations
1846            .contains_key(&capability_key)
1847        {
1848            return Err(format!(
1849                "native capability {capability_key} is not registered"
1850            ));
1851        }
1852    }
1853    // The external desired-peer actor is one Rust root actor. Public avenue
1854    // principals are registration metadata, not lifecycle keys.
1855    let root_user_id = format!("{}:native-root", state.client().app_tag());
1856    let _requested_capability_principal = user_id;
1857    state
1858        .client()
1859        .start_external_auto_connect(root_user_id, local_device_id)
1860        .await
1861        .map_err(|error| error.to_string())
1862}
1863
1864#[tauri::command]
1865async fn submit_rtc_desired_peers(
1866    state: tauri::State<'_, OpenRtcTauriState>,
1867    revision: u64,
1868    peers_json: String,
1869    capability_key: Option<String>,
1870) -> Result<bool, String> {
1871    let peers = serde_json::from_str::<Vec<serde_json::Value>>(&peers_json)
1872        .map_err(|error| format!("invalid capability desired-peer payload: {error}"))?;
1873    if peers.len() > 100 {
1874        return Err("capability desired-peer payload exceeds 100 peers".to_string());
1875    }
1876    let (root_revision, aggregate) = {
1877        let mut registry = state.capability_registry.lock().await;
1878        let key = match capability_key {
1879            Some(key) => required_native_capability_part(key, "capability key")?,
1880            None if registry.registrations.len() == 1 => registry
1881                .registrations
1882                .keys()
1883                .next()
1884                .cloned()
1885                .expect("one registration has one key"),
1886            None => {
1887                return Err(
1888                    "capabilityKey is required when multiple native capabilities are active"
1889                        .to_string(),
1890                )
1891            }
1892        };
1893        let registration = registry
1894            .registrations
1895            .get_mut(&key)
1896            .ok_or_else(|| format!("native capability {key} is not registered"))?;
1897        if revision <= registration.desired_revision {
1898            return Ok(false);
1899        }
1900        registration.desired_revision = revision;
1901        registration.desired_peers = peers;
1902        registry.aggregate_desired_peers()?
1903    };
1904    let accepted = state
1905        .client()
1906        .submit_external_desired_peers(root_revision, &aggregate)
1907        .await
1908        .map_err(|error| error.to_string())?;
1909    if accepted {
1910        // A capability can be registered after its peer is already connected.
1911        // Replay the Rust-owned read model after the desired-peer mapping is
1912        // committed instead of asking the frontend to redial or poll.
1913        replay_current_connection_states(&state).await;
1914    }
1915    Ok(accepted)
1916}
1917
1918#[tauri::command]
1919async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1920    state.client().stop_auto_connect();
1921    state.client().stop_external_auto_connect().await;
1922    if let Some(active) = state.managed_session.lock().await.as_mut() {
1923        active.result.auto_connect_started = false;
1924    }
1925    Ok(())
1926}
1927
1928#[tauri::command]
1929async fn connect_to_device(
1930    state: tauri::State<'_, OpenRtcTauriState>,
1931    device_id: Option<String>,
1932    endpoint_ticket: String,
1933    timeout_ms: Option<u64>,
1934) -> Result<openrtc::client::ManagedConnectResult, String> {
1935    let client = state.client();
1936    let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
1937    match timeout_ms {
1938        Some(timeout_ms) => {
1939            tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
1940                .await
1941                .map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
1942                .map_err(|error| error.to_string())
1943        }
1944        None => connect.await.map_err(|error| error.to_string()),
1945    }
1946}
1947
1948#[tauri::command]
1949async fn disconnect_device(
1950    state: tauri::State<'_, OpenRtcTauriState>,
1951    device_id: String,
1952    node_id_hint: Option<String>,
1953) -> Result<(), String> {
1954    state
1955        .client()
1956        .disconnect_device(&device_id, node_id_hint.as_deref())
1957        .await;
1958    Ok(())
1959}
1960
1961#[tauri::command]
1962async fn set_auto_connect_excluded(
1963    state: tauri::State<'_, OpenRtcTauriState>,
1964    device_id: String,
1965    excluded: bool,
1966) -> Result<(), String> {
1967    if excluded {
1968        state.client().exclude_peer_and_publish(&device_id).await;
1969    } else {
1970        state.client().unexclude_peer_and_publish(&device_id).await;
1971    }
1972    Ok(())
1973}
1974
1975#[tauri::command]
1976async fn set_rtc_external_auto_connect_excluded(
1977    state: tauri::State<'_, OpenRtcTauriState>,
1978    device_id: String,
1979    excluded: bool,
1980) -> Result<(), String> {
1981    state
1982        .client()
1983        .set_auto_connect_excluded(&device_id, excluded);
1984    state.client().wake_native_external_auto_connect().await;
1985    Ok(())
1986}
1987
1988#[tauri::command]
1989async fn resolve_rtc_peer_connection_records(
1990    state: tauri::State<'_, OpenRtcTauriState>,
1991    id: String,
1992) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1993    Ok(state.client().resolve_peer_connection_records(&id).await)
1994}
1995
1996#[tauri::command]
1997async fn resolve_rtc_peer_identity(
1998    state: tauri::State<'_, OpenRtcTauriState>,
1999    id: String,
2000) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
2001    Ok(state.client().peer_snapshot(&id).await)
2002}
2003
2004#[tauri::command]
2005async fn get_rtc_peer_session(
2006    state: tauri::State<'_, OpenRtcTauriState>,
2007    id: String,
2008) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2009    Ok(state.client().peer_session(&id).await)
2010}
2011
2012#[tauri::command]
2013async fn list_rtc_peer_sessions(
2014    state: tauri::State<'_, OpenRtcTauriState>,
2015) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
2016    Ok(state.client().peer_sessions().await)
2017}
2018
2019#[tauri::command]
2020async fn list_rtc_managed_connections(
2021    state: tauri::State<'_, OpenRtcTauriState>,
2022) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
2023    Ok(state.client().list_managed_connections().await)
2024}
2025
2026#[derive(Debug, Clone, Serialize)]
2027#[serde(rename_all = "camelCase")]
2028struct NotifyRtcNetworkChangeResult {
2029    retired_stale_connections: usize,
2030}
2031
2032#[tauri::command]
2033async fn notify_rtc_network_change(
2034    state: tauri::State<'_, OpenRtcTauriState>,
2035) -> Result<NotifyRtcNetworkChangeResult, String> {
2036    let retired_stale_connections = state
2037        .client()
2038        .notify_network_change()
2039        .await
2040        .map_err(|error| error.to_string())?;
2041    Ok(NotifyRtcNetworkChangeResult {
2042        retired_stale_connections,
2043    })
2044}
2045
2046#[tauri::command]
2047async fn set_rtc_transport_priority(
2048    state: tauri::State<'_, OpenRtcTauriState>,
2049    priority: Vec<openrtc::route_policy::KnownRoute>,
2050) -> Result<(), String> {
2051    let client = state.client();
2052    let mut transport_config = client.transport_config().await;
2053    transport_config.route_priority = priority;
2054    client
2055        .update_transport_config(transport_config)
2056        .await
2057        .map_err(|error| error.to_string())
2058}
2059
2060#[tauri::command]
2061async fn wait_for_rtc_settled_peer(
2062    state: tauri::State<'_, OpenRtcTauriState>,
2063    id: String,
2064    timeout_ms: Option<u64>,
2065) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2066    Ok(state.client().wait_for_peer(&id, timeout_ms).await)
2067}
2068
2069#[tauri::command]
2070async fn list_rtc_connection_states(
2071    state: tauri::State<'_, OpenRtcTauriState>,
2072    capability_key: String,
2073) -> Result<Vec<openrtc::client::StateSnapshot>, String> {
2074    let states = state.client().connection_states().await;
2075    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2076    let registry = state.capability_registry.lock().await;
2077    Ok(states
2078        .into_iter()
2079        .filter(|snapshot| {
2080            registry
2081                .capability_keys_for_state(snapshot)
2082                .iter()
2083                .any(|key| key == &capability_key)
2084        })
2085        .collect())
2086}
2087
2088#[tauri::command]
2089async fn get_rtc_connection_state(
2090    state: tauri::State<'_, OpenRtcTauriState>,
2091    connection_id: String,
2092    capability_key: String,
2093) -> Result<Option<openrtc::client::StateSnapshot>, String> {
2094    let snapshot = state.client().connection_state(&connection_id).await;
2095    let Some(snapshot) = snapshot else {
2096        return Ok(None);
2097    };
2098    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2099    let included = state
2100        .capability_registry
2101        .lock()
2102        .await
2103        .capability_keys_for_state(&snapshot)
2104        .iter()
2105        .any(|key| key == &capability_key);
2106    Ok(included.then_some(snapshot))
2107}
2108
2109#[tauri::command]
2110async fn stop_rtc_subscription(
2111    state: tauri::State<'_, OpenRtcTauriState>,
2112    request_id: String,
2113) -> Result<(), String> {
2114    stop_subscription_by_id(&state, &request_id).await;
2115    Ok(())
2116}
2117
2118async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
2119    if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
2120        handle.abort();
2121    }
2122    let mut projected = state.native_projection_subscription.lock().await;
2123    if projected
2124        .as_ref()
2125        .is_some_and(|subscription| subscription.request_id == request_id)
2126    {
2127        *projected = None;
2128    }
2129}
2130
2131/// Register the one capability-keyed WebView projection Channel.
2132///
2133/// Rust forwards canonical connection state, decrypted peer data, and streams
2134/// that the embedding host explicitly hands to `handoff_projected_peer_bi_stream`.
2135/// The Channel is projection only; it never attaches another Rust receiver or
2136/// owns peer lifecycle.
2137#[tauri::command]
2138async fn start_native_projection(
2139    state: tauri::State<'_, OpenRtcTauriState>,
2140    request_id: Option<String>,
2141    channel: Channel<NativeProjectionEvent>,
2142) -> Result<String, String> {
2143    let request_id = request_id
2144        .map(|value| value.trim().to_string())
2145        .filter(|value| !value.is_empty())
2146        .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
2147    let mut subscription = state.native_projection_subscription.lock().await;
2148    if let Some(current) = subscription.as_ref() {
2149        if current.request_id == request_id {
2150            return Ok(request_id);
2151        }
2152        return Err(
2153            "OpenRTC native projection subscription already has a process owner".to_string(),
2154        );
2155    }
2156    *subscription = Some(NativeProjectionSubscription {
2157        request_id: request_id.clone(),
2158        channel,
2159    });
2160    Ok(request_id)
2161}
2162
2163#[tauri::command]
2164async fn get_projected_stream_diagnostics(
2165    state: tauri::State<'_, OpenRtcTauriState>,
2166) -> Result<ProjectedStreamDiagnostics, String> {
2167    Ok(ProjectedStreamDiagnostics {
2168        subscription_active: state.native_projection_subscription.lock().await.is_some(),
2169        offered: state.projected_stream_offered.load(Ordering::Relaxed),
2170        decoded: state.projected_stream_decoded.load(Ordering::Relaxed),
2171        projected: state.projected_stream_projected.load(Ordering::Relaxed),
2172        unhandled_non_channel: state
2173            .projected_stream_unhandled_non_channel
2174            .load(Ordering::Relaxed),
2175        unhandled_other_channel: state
2176            .projected_stream_unhandled_other_channel
2177            .load(Ordering::Relaxed),
2178        unhandled_no_subscription: state
2179            .projected_stream_unhandled_no_subscription
2180            .load(Ordering::Relaxed),
2181        failures: state.projected_stream_failures.load(Ordering::Relaxed),
2182        last_transport_stable_id: state
2183            .projected_stream_last_transport_stable_id
2184            .load(Ordering::Relaxed),
2185        validation_checks: state
2186            .projected_stream_validation_checks
2187            .load(Ordering::Relaxed),
2188        validation_rejections: state
2189            .projected_stream_validation_rejections
2190            .load(Ordering::Relaxed),
2191        last_validated_transport_stable_id: state
2192            .projected_stream_last_validated_transport_stable_id
2193            .load(Ordering::Relaxed),
2194        authorized: state.projected_stream_authorized.load(Ordering::Relaxed),
2195        unauthorized: state.projected_stream_unauthorized.load(Ordering::Relaxed),
2196        last_channel: state.projected_stream_last_channel.lock().await.clone(),
2197        last_protocol: state.projected_stream_last_protocol.lock().await.clone(),
2198        last_connection_id: state
2199            .projected_stream_last_connection_id
2200            .lock()
2201            .await
2202            .clone(),
2203        last_remote_node_id: state
2204            .projected_stream_last_remote_node_id
2205            .lock()
2206            .await
2207            .clone(),
2208    })
2209}
2210
2211#[tauri::command]
2212async fn is_current_transport_stable_id(
2213    state: tauri::State<'_, OpenRtcTauriState>,
2214    endpoint_id: String,
2215    transport_stable_id: u64,
2216) -> Result<bool, String> {
2217    state
2218        .projected_stream_validation_checks
2219        .fetch_add(1, Ordering::Relaxed);
2220    state
2221        .projected_stream_last_validated_transport_stable_id
2222        .store(transport_stable_id, Ordering::Relaxed);
2223    let is_current = state
2224        .client()
2225        .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
2226        .await
2227        .map_err(|error| error.to_string())?;
2228    if !is_current {
2229        state
2230            .projected_stream_validation_rejections
2231            .fetch_add(1, Ordering::Relaxed);
2232    }
2233    if native_stream_trace_enabled() {
2234        eprintln!(
2235            "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
2236            endpoint_id, transport_stable_id, is_current
2237        );
2238    }
2239    Ok(is_current)
2240}
2241
2242#[tauri::command]
2243async fn send_peer_message(
2244    state: tauri::State<'_, OpenRtcTauriState>,
2245    request: Request<'_>,
2246) -> Result<(), String> {
2247    let id = required_request_header(&request, PEER_ID_HEADER)?;
2248    let data = request_body_bytes(&request)?;
2249    state
2250        .client()
2251        .send_peer(&id, &data)
2252        .await
2253        .map_err(|error| error.to_string())
2254}
2255
2256#[tauri::command]
2257async fn is_peer_connected(
2258    state: tauri::State<'_, OpenRtcTauriState>,
2259    node_id: String,
2260) -> Result<bool, String> {
2261    state
2262        .client()
2263        .is_connected_str(&node_id)
2264        .await
2265        .map_err(|error| error.to_string())
2266}
2267
2268#[derive(Debug, Serialize)]
2269#[serde(rename_all = "camelCase")]
2270struct PortableMediaPublicationStart {
2271    publication_id: String,
2272    media_generation: u32,
2273    control: Vec<u8>,
2274}
2275
2276fn media_kind_name(kind: openrtc::media::MediaKind) -> &'static str {
2277    match kind {
2278        openrtc::media::MediaKind::Audio => "audio",
2279        openrtc::media::MediaKind::Video => "video",
2280        openrtc::media::MediaKind::Screen => "screen",
2281        openrtc::media::MediaKind::Data => "data",
2282    }
2283}
2284
2285fn media_codec_name(codec: openrtc::media::MediaCodec) -> &'static str {
2286    match codec {
2287        openrtc::media::MediaCodec::Opus => "opus",
2288        openrtc::media::MediaCodec::H264 => "h264",
2289        openrtc::media::MediaCodec::Vp8 => "vp8",
2290        openrtc::media::MediaCodec::Vp9 => "vp9",
2291        openrtc::media::MediaCodec::Av1 => "av1",
2292        openrtc::media::MediaCodec::Pcm => "pcm",
2293        openrtc::media::MediaCodec::Opaque => "opaque",
2294    }
2295}
2296
2297fn parse_media_publication_id(value: &str) -> Result<openrtc::media::PublicationId, String> {
2298    value
2299        .trim()
2300        .parse()
2301        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())
2302}
2303
2304#[tauri::command]
2305#[allow(clippy::too_many_arguments)]
2306async fn begin_openrtc_media_publication(
2307    state: tauri::State<'_, OpenRtcTauriState>,
2308    publication_id: Option<String>,
2309    kind: String,
2310    codec: String,
2311    clock_rate: u32,
2312    coded_width: Option<u32>,
2313    coded_height: Option<u32>,
2314    channels: Option<u16>,
2315) -> Result<PortableMediaPublicationStart, String> {
2316    let publication_id = match publication_id.as_deref().map(str::trim) {
2317        Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
2318        _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
2319    };
2320    let kind: openrtc::media::MediaKind = kind
2321        .parse()
2322        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
2323    if !matches!(
2324        kind,
2325        openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
2326    ) {
2327        return Err("browser media publications must be audio or video".to_string());
2328    }
2329    let codec = codec
2330        .parse()
2331        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
2332    let publication = openrtc::media::MediaPublicationConfig {
2333        publication_id,
2334        media_generation: 0,
2335        kind,
2336        codec,
2337        clock_rate,
2338        coded_width,
2339        coded_height,
2340        channels,
2341    };
2342    let (publication, control) = state
2343        .portable_media
2344        .lock()
2345        .await
2346        .begin_publication(publication)
2347        .map_err(|error| error.to_string())?;
2348    Ok(PortableMediaPublicationStart {
2349        publication_id: publication.publication_id.to_string(),
2350        media_generation: publication.media_generation,
2351        control,
2352    })
2353}
2354
2355#[tauri::command]
2356#[allow(clippy::too_many_arguments)]
2357async fn encode_openrtc_media_sample(
2358    state: tauri::State<'_, OpenRtcTauriState>,
2359    request: Request<'_>,
2360) -> Result<Response, String> {
2361    let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
2362    let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
2363    openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
2364        .map_err(|error| error.to_string())?;
2365    let duration_us = parsed_request_header(&request, MEDIA_DURATION_US_HEADER)?;
2366    let keyframe = parsed_request_header(&request, MEDIA_KEYFRAME_HEADER)?;
2367    let discardable = parsed_request_header(&request, MEDIA_DISCARDABLE_HEADER)?;
2368    let payload = request_body_bytes(&request)?;
2369    let encoded = state
2370        .portable_media
2371        .lock()
2372        .await
2373        .encode_sample(
2374            parse_media_publication_id(&publication_id)?,
2375            openrtc::media::EncodedMediaSample {
2376                timestamp_us,
2377                duration_us,
2378                keyframe,
2379                discardable,
2380                payload,
2381            },
2382        )
2383        .map_err(|error| error.to_string())?;
2384    Ok(Response::new(encoded))
2385}
2386
2387#[tauri::command]
2388async fn pause_openrtc_media_publication(
2389    state: tauri::State<'_, OpenRtcTauriState>,
2390    publication_id: String,
2391) -> Result<(), String> {
2392    state
2393        .portable_media
2394        .lock()
2395        .await
2396        .pause_publication(parse_media_publication_id(&publication_id)?)
2397        .map_err(|error| error.to_string())
2398}
2399
2400#[tauri::command]
2401async fn retire_openrtc_media_publication(
2402    state: tauri::State<'_, OpenRtcTauriState>,
2403    publication_id: String,
2404) -> Result<(), String> {
2405    state
2406        .portable_media
2407        .lock()
2408        .await
2409        .retire_publication(parse_media_publication_id(&publication_id)?);
2410    Ok(())
2411}
2412
2413#[tauri::command]
2414async fn retire_openrtc_media_receiver(
2415    state: tauri::State<'_, OpenRtcTauriState>,
2416    publication_id: String,
2417    media_generation: u32,
2418) -> Result<(), String> {
2419    state.portable_media.lock().await.retire_receiver(
2420        parse_media_publication_id(&publication_id)?,
2421        media_generation,
2422    );
2423    Ok(())
2424}
2425
2426#[tauri::command]
2427async fn decode_openrtc_media_chunk(
2428    state: tauri::State<'_, OpenRtcTauriState>,
2429    request: Request<'_>,
2430) -> Result<Response, String> {
2431    let encoded = request_body_bytes(&request)?;
2432    let chunk = state
2433        .portable_media
2434        .lock()
2435        .await
2436        .decode_chunk(&encoded)
2437        .map_err(|error| error.to_string())?;
2438    openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
2439        .and_then(|_| {
2440            openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
2441        })
2442        .map_err(|error| error.to_string())?;
2443    let metadata = serde_json::to_vec(&serde_json::json!({
2444        "publicationId": chunk.publication_id.to_string(),
2445        "mediaGeneration": chunk.media_generation,
2446        "sequence": chunk.sequence,
2447        "timestampUs": chunk.timestamp_us,
2448        "durationUs": chunk.duration_us,
2449        "kind": media_kind_name(chunk.kind),
2450        "codec": media_codec_name(chunk.codec),
2451        "keyframe": chunk.keyframe,
2452        "discardable": chunk.discardable,
2453    }))
2454    .map_err(|error| format!("encode media IPC metadata failed: {error}"))?;
2455    let metadata_len =
2456        u32::try_from(metadata.len()).map_err(|_| "media IPC metadata is too large".to_string())?;
2457    let mut response = Vec::with_capacity(4 + metadata.len() + chunk.payload.len());
2458    response.extend_from_slice(&metadata_len.to_be_bytes());
2459    response.extend_from_slice(&metadata);
2460    response.extend_from_slice(&chunk.payload);
2461    Ok(Response::new(response))
2462}
2463
2464#[tauri::command]
2465async fn decode_openrtc_media_control(
2466    state: tauri::State<'_, OpenRtcTauriState>,
2467    request: Request<'_>,
2468) -> Result<serde_json::Value, String> {
2469    let encoded = request_body_bytes(&request)?;
2470    let control = state
2471        .portable_media
2472        .lock()
2473        .await
2474        .decode_control(&encoded)
2475        .map_err(|error| error.to_string())?;
2476    Ok(match control {
2477        openrtc::media::MediaControlFrame::Publish {
2478            publication_id,
2479            media_generation,
2480            kind,
2481            codec,
2482            clock_rate,
2483            coded_width,
2484            coded_height,
2485            channels,
2486        } => serde_json::json!({
2487            "type": "publish",
2488            "publicationId": publication_id.to_string(),
2489            "mediaGeneration": media_generation,
2490            "kind": media_kind_name(kind),
2491            "codec": media_codec_name(codec),
2492            "clockRate": clock_rate,
2493            "codedWidth": coded_width,
2494            "codedHeight": coded_height,
2495            "channels": channels,
2496        }),
2497        openrtc::media::MediaControlFrame::SetEnabled {
2498            publication_id,
2499            media_generation,
2500            enabled,
2501        } => serde_json::json!({
2502            "type": "set-enabled",
2503            "publicationId": publication_id.to_string(),
2504            "mediaGeneration": media_generation,
2505            "enabled": enabled,
2506        }),
2507        openrtc::media::MediaControlFrame::RequestKeyframe {
2508            publication_id,
2509            media_generation,
2510        } => serde_json::json!({
2511            "type": "request-keyframe",
2512            "publicationId": publication_id.to_string(),
2513            "mediaGeneration": media_generation,
2514        }),
2515        openrtc::media::MediaControlFrame::Stop {
2516            publication_id,
2517            media_generation,
2518            reason,
2519        } => serde_json::json!({
2520            "type": "stop",
2521            "publicationId": publication_id.to_string(),
2522            "mediaGeneration": media_generation,
2523            "reason": reason,
2524        }),
2525    })
2526}
2527
2528#[tauri::command]
2529async fn open_peer_bi_stream<R: Runtime>(
2530    app: tauri::AppHandle<R>,
2531    state: tauri::State<'_, OpenRtcTauriState>,
2532    peer_id: String,
2533    timeout_ms: Option<u64>,
2534) -> Result<OpenBiResult, String> {
2535    open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
2536        client.open_peer_bi(&peer_id, timeout_ms).await
2537    })
2538    .await
2539}
2540
2541#[tauri::command]
2542async fn open_peer_bi_transport_only_stream<R: Runtime>(
2543    app: tauri::AppHandle<R>,
2544    state: tauri::State<'_, OpenRtcTauriState>,
2545    peer_id: String,
2546    timeout_ms: Option<u64>,
2547) -> Result<OpenBiResult, String> {
2548    open_peer_bi_with(
2549        app,
2550        state,
2551        peer_id,
2552        timeout_ms,
2553        |client, peer_id, timeout_ms| async move {
2554            client
2555                .open_peer_bi_transport_only(&peer_id, timeout_ms)
2556                .await
2557                .map(|(connection_id, remote_node_id, send, recv)| {
2558                    (
2559                        connection_id,
2560                        remote_node_id,
2561                        PeerSendStream::plain(send),
2562                        PeerRecvStream::plain(recv),
2563                    )
2564                })
2565        },
2566    )
2567    .await
2568}
2569
2570#[tauri::command]
2571async fn open_peer_uni_stream(
2572    state: tauri::State<'_, OpenRtcTauriState>,
2573    peer_id: String,
2574    timeout_ms: Option<u64>,
2575) -> Result<OpenUniResult, String> {
2576    let peer_id = peer_id.trim().to_string();
2577    if peer_id.is_empty() {
2578        return Err("peerId is required".to_string());
2579    }
2580    let (connection_id, remote_node_id, send) = state
2581        .client()
2582        .open_peer_uni(&peer_id, timeout_ms)
2583        .await
2584        .map_err(|error| format!("open peer uni stream failed: {error}"))?;
2585    let stream_id = uuid::Uuid::new_v4().to_string();
2586    state
2587        .peer_uni_streams
2588        .lock()
2589        .await
2590        .insert(stream_id.clone(), Arc::new(Mutex::new(Some(send))));
2591    Ok(OpenUniResult {
2592        stream_id,
2593        connection_id,
2594        remote_node_id,
2595    })
2596}
2597
2598#[tauri::command]
2599async fn write_peer_bi_stream(
2600    state: tauri::State<'_, OpenRtcTauriState>,
2601    request: Request<'_>,
2602) -> Result<(), String> {
2603    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
2604    let bytes = request_body_bytes(&request)?;
2605
2606    let send = {
2607        let streams = state.peer_bi_streams.lock().await;
2608        streams
2609            .get(&stream_id)
2610            .map(|handle| handle.send.clone())
2611            .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
2612    };
2613    let result = send
2614        .lock()
2615        .await
2616        .as_mut()
2617        .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
2618        .write_all(&bytes)
2619        .await
2620        .map_err(|error| format!("write peer bi stream failed: {error}"));
2621    result
2622}
2623
2624#[tauri::command]
2625async fn start_peer_bi_stream_read<R: Runtime>(
2626    _app: tauri::AppHandle<R>,
2627    state: tauri::State<'_, OpenRtcTauriState>,
2628    stream_id: String,
2629    channel: Channel<Response>,
2630) -> Result<(), String> {
2631    let stream_id = stream_id.trim().to_string();
2632    if stream_id.is_empty() {
2633        return Err("streamId is required".to_string());
2634    }
2635
2636    let mut streams = state.peer_bi_streams.lock().await;
2637    let handle = streams
2638        .get_mut(&stream_id)
2639        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
2640    if handle.read_task.is_some() {
2641        return Ok(());
2642    }
2643    let recv = handle
2644        .recv
2645        .take()
2646        .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
2647    handle.read_task = Some(spawn_peer_bi_stream_reader(channel, stream_id, recv));
2648    Ok(())
2649}
2650
2651/// Cancel only the receive half of a bidirectional stream.
2652///
2653/// Web Streams permits a consumer to cancel an unused readable while keeping
2654/// its writable alive. Media publications rely on that half-close behavior, so
2655/// the native IPC projection must not remove or finish the Rust send handle.
2656#[tauri::command]
2657async fn cancel_peer_bi_stream_read(
2658    state: tauri::State<'_, OpenRtcTauriState>,
2659    stream_id: String,
2660) -> Result<(), String> {
2661    let stream_id = stream_id.trim().to_string();
2662    if stream_id.is_empty() {
2663        return Err("streamId is required".to_string());
2664    }
2665
2666    let mut streams = state.peer_bi_streams.lock().await;
2667    let handle = streams
2668        .get_mut(&stream_id)
2669        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
2670    handle.recv.take();
2671    if let Some(read_task) = handle.read_task.take() {
2672        read_task.abort();
2673    }
2674    Ok(())
2675}
2676
2677#[tauri::command]
2678async fn close_peer_bi_stream(
2679    state: tauri::State<'_, OpenRtcTauriState>,
2680    stream_id: String,
2681) -> Result<(), String> {
2682    let stream_id = stream_id.trim().to_string();
2683    if stream_id.is_empty() {
2684        return Err("streamId is required".to_string());
2685    }
2686
2687    let handle = {
2688        let mut streams = state.peer_bi_streams.lock().await;
2689        streams.remove(&stream_id)
2690    };
2691    if let Some(handle) = handle {
2692        // Hold the writer mutex before finishing so close waits for any
2693        // in-flight chunk write instead of skipping the flush when Arc clones exist.
2694        let mut send = handle.send.lock().await;
2695        if let Some(send) = send.take() {
2696            let _ = send.finish();
2697        }
2698        drop(send);
2699        if let Some(read_task) = handle.read_task {
2700            read_task.abort();
2701        }
2702    }
2703    Ok(())
2704}
2705
2706#[tauri::command]
2707async fn write_peer_uni_stream(
2708    state: tauri::State<'_, OpenRtcTauriState>,
2709    request: Request<'_>,
2710) -> Result<(), String> {
2711    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
2712    let bytes = request_body_bytes(&request)?;
2713    let send = {
2714        let streams = state.peer_uni_streams.lock().await;
2715        streams
2716            .get(&stream_id)
2717            .cloned()
2718            .ok_or_else(|| format!("peer uni stream not found: {stream_id}"))?
2719    };
2720    let result = send
2721        .lock()
2722        .await
2723        .as_mut()
2724        .ok_or_else(|| format!("peer uni stream already closed: {stream_id}"))?
2725        .write_all(&bytes)
2726        .await
2727        .map_err(|error| format!("write peer uni stream failed: {error}"));
2728    result
2729}
2730
2731#[tauri::command]
2732async fn close_peer_uni_stream(
2733    state: tauri::State<'_, OpenRtcTauriState>,
2734    stream_id: String,
2735) -> Result<(), String> {
2736    let stream_id = stream_id.trim().to_string();
2737    if stream_id.is_empty() {
2738        return Err("streamId is required".to_string());
2739    }
2740    let send = state.peer_uni_streams.lock().await.remove(&stream_id);
2741    if let Some(send) = send {
2742        if let Some(send) = send.lock().await.take() {
2743            let _ = send.finish();
2744        }
2745    }
2746    Ok(())
2747}
2748
2749#[cfg(test)]
2750mod tests {
2751    use super::*;
2752    use std::collections::HashMap as StdHashMap;
2753    use std::sync::atomic::{AtomicUsize, Ordering};
2754    use std::sync::Mutex as StdMutex;
2755    use tokio::sync::oneshot;
2756
2757    struct FakeTransportInstaller {
2758        installs: Arc<AtomicUsize>,
2759    }
2760
2761    #[derive(Default)]
2762    struct FakeDeviceKeySigner {
2763        records: StdMutex<StdHashMap<String, String>>,
2764    }
2765
2766    impl DeviceKeySigner for FakeDeviceKeySigner {
2767        fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
2768            Ok(serde_json::json!({
2769                "kty": "OKP",
2770                "crv": "Ed25519",
2771                "x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
2772            }))
2773        }
2774
2775        fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
2776            assert_eq!(challenge, b"openrtc:v2:test");
2777            Ok(vec![7; 64])
2778        }
2779
2780        fn delete(&self, _app_tag: &str) -> Result<(), String> {
2781            Ok(())
2782        }
2783
2784        fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
2785            Ok(self
2786                .records
2787                .lock()
2788                .unwrap()
2789                .get(&format!("{app_tag}:{key}"))
2790                .cloned())
2791        }
2792
2793        fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
2794            self.records
2795                .lock()
2796                .unwrap()
2797                .insert(format!("{app_tag}:{key}"), value.to_string());
2798            Ok(())
2799        }
2800
2801        fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
2802            self.records
2803                .lock()
2804                .unwrap()
2805                .remove(&format!("{app_tag}:{key}"));
2806            Ok(())
2807        }
2808    }
2809
2810    impl TransportInstaller for FakeTransportInstaller {
2811        fn id(&self) -> &'static str {
2812            "ble"
2813        }
2814
2815        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2816            config.ble.as_ref().is_some_and(|ble| ble.enabled)
2817        }
2818
2819        fn install(
2820            &self,
2821            _client: Arc<openrtc::client::Client>,
2822            _config: openrtc::client::TransportConfig,
2823            _context: InstallContext,
2824        ) -> InstallFuture {
2825            self.installs.fetch_add(1, Ordering::SeqCst);
2826            Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
2827        }
2828    }
2829
2830    #[derive(Debug)]
2831    struct UnavailableTransportInstaller;
2832
2833    impl TransportInstaller for UnavailableTransportInstaller {
2834        fn id(&self) -> &'static str {
2835            "ble"
2836        }
2837
2838        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2839            config.ble.as_ref().is_some_and(|ble| ble.enabled)
2840        }
2841
2842        fn install(
2843            &self,
2844            _client: Arc<openrtc::client::Client>,
2845            _config: openrtc::client::TransportConfig,
2846            _context: InstallContext,
2847        ) -> InstallFuture {
2848            Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
2849        }
2850    }
2851
2852    fn requested_ble_config() -> openrtc::client::TransportConfig {
2853        openrtc::client::TransportConfig {
2854            ble: Some(openrtc::client::BleConfig {
2855                enabled: true,
2856                ..openrtc::client::BleConfig::default()
2857            }),
2858            ..openrtc::client::TransportConfig::default()
2859        }
2860    }
2861
2862    fn test_config() -> OpenRtcTauriConfig {
2863        OpenRtcTauriConfig {
2864            api_key: format!("pk_test_{}", "a".repeat(40)),
2865            ..OpenRtcTauriConfig::default()
2866        }
2867    }
2868
2869    #[test]
2870    fn native_capability_registry_isolates_duplicate_peer_channels() {
2871        let mut registry = NativeCapabilityRegistry::default();
2872        registry
2873            .register(
2874                "devices:user-1".to_string(),
2875                "devices".to_string(),
2876                "user-1".to_string(),
2877            )
2878            .expect("devices registration");
2879        registry
2880            .register(
2881                "space:room-1".to_string(),
2882                "space".to_string(),
2883                "room-1".to_string(),
2884            )
2885            .expect("space registration");
2886        registry
2887            .registrations
2888            .get_mut("devices:user-1")
2889            .unwrap()
2890            .desired_peers = vec![serde_json::json!({
2891            "deviceId": "shared-peer",
2892            "nodeId": "device-route"
2893        })];
2894        registry
2895            .registrations
2896            .get_mut("space:room-1")
2897            .unwrap()
2898            .desired_peers = vec![serde_json::json!({
2899            "deviceId": "shared-peer",
2900            "nodeId": "space-route"
2901        })];
2902
2903        assert_eq!(
2904            registry.capability_keys_for_identity(
2905                Some("shared-peer"),
2906                Some("shared-peer"),
2907                None,
2908                None,
2909            ),
2910            vec!["devices:user-1".to_string(), "space:room-1".to_string()]
2911        );
2912
2913        let explicit_channel = openrtc::stream_metadata::ChannelMetadata {
2914            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
2915            metadata: Some(serde_json::Map::from_iter([(
2916                "openrtcCapability".to_string(),
2917                serde_json::json!("space:room-1"),
2918            )])),
2919        };
2920        assert_eq!(
2921            registry.projected_stream_capability(&explicit_channel),
2922            Some("space:room-1".to_string())
2923        );
2924
2925        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
2926            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
2927            metadata: None,
2928        };
2929        assert_eq!(
2930            registry.projected_stream_capability(&unscoped_channel),
2931            None
2932        );
2933    }
2934
2935    #[test]
2936    fn native_capability_registry_aggregates_without_duplicating_root_peers() {
2937        let mut registry = NativeCapabilityRegistry::default();
2938        for (key, kind, id) in [
2939            ("devices:user-1", "devices", "user-1"),
2940            ("space:room-1", "space", "room-1"),
2941        ] {
2942            registry
2943                .register(key.to_string(), kind.to_string(), id.to_string())
2944                .expect("capability registration");
2945            registry.registrations.get_mut(key).unwrap().desired_peers =
2946                vec![serde_json::json!({"deviceId": "shared-peer", "ticket": "ticket-a"})];
2947        }
2948
2949        let (first_revision, peers_json) = registry.aggregate_desired_peers().unwrap();
2950        let peers: Vec<serde_json::Value> = serde_json::from_str(&peers_json).unwrap();
2951        assert_eq!(first_revision, 1);
2952        assert_eq!(peers.len(), 1, "one physical peer must be dialed once");
2953
2954        registry
2955            .register(
2956                "devices:user-1".to_string(),
2957                "devices".to_string(),
2958                "user-1".to_string(),
2959            )
2960            .expect("exact registration is idempotent");
2961        assert!(registry
2962            .register(
2963                "devices:user-1".to_string(),
2964                "space".to_string(),
2965                "room-1".to_string(),
2966            )
2967            .is_err());
2968    }
2969
2970    #[test]
2971    fn native_capability_registry_scopes_peer_data_and_rejects_unscoped_streams() {
2972        let mut registry = NativeCapabilityRegistry::default();
2973        registry
2974            .register(
2975                "devices:user-1".to_string(),
2976                "devices".to_string(),
2977                "user-1".to_string(),
2978            )
2979            .unwrap();
2980        registry
2981            .registrations
2982            .get_mut("devices:user-1")
2983            .unwrap()
2984            .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
2985
2986        let explicit = openrtc::client::NativePeerDataEvent {
2987            connection_id: "peer-1".to_string(),
2988            remote_node_id: None,
2989            transport: "webrtc".to_string(),
2990            transport_stable_id: 1,
2991            transport_generation: 1,
2992            route_generation: 1,
2993            payload: serde_json::to_vec(&serde_json::json!({
2994                "capability": "devices:user-1",
2995                "kind": "raw",
2996                "payload": [1, 2, 3]
2997            }))
2998            .unwrap(),
2999        };
3000        assert_eq!(
3001            registry.capability_keys_for_peer_data(&explicit),
3002            vec!["devices:user-1".to_string()]
3003        );
3004
3005        let unknown = openrtc::client::NativePeerDataEvent {
3006            payload: serde_json::to_vec(&serde_json::json!({
3007                "capability": "space:unknown",
3008                "kind": "raw"
3009            }))
3010            .unwrap(),
3011            ..explicit.clone()
3012        };
3013        assert!(registry.capability_keys_for_peer_data(&unknown).is_empty());
3014
3015        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
3016            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
3017            metadata: None,
3018        };
3019        assert_eq!(
3020            registry.projected_stream_capability(&unscoped_channel),
3021            None,
3022            "stream ownership must never be inferred from capability count"
3023        );
3024    }
3025
3026    #[test]
3027    fn native_device_proof_uses_host_signer_without_exporting_private_key() {
3028        let state = OpenRtcTauriState::new(test_config())
3029            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
3030        let public = state
3031            .device_public_key("app_native_test")
3032            .expect("public key");
3033        assert_eq!(public["kty"], "OKP");
3034        assert_eq!(public["crv"], "Ed25519");
3035        assert!(public.get("d").is_none());
3036        let signature = state
3037            .sign_device_proof("app_native_test", "openrtc:v2:test")
3038            .expect("signature");
3039        assert_eq!(
3040            base64::engine::general_purpose::URL_SAFE_NO_PAD
3041                .decode(signature)
3042                .expect("base64"),
3043            vec![7; 64]
3044        );
3045    }
3046
3047    #[test]
3048    fn native_certificate_records_round_trip_through_host_secure_storage() {
3049        let state = OpenRtcTauriState::new(test_config())
3050            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
3051        let app_tag = "app_native_test";
3052        let key = "openrtc:v2:device-session:abc:device-1";
3053        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
3054        state
3055            .write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
3056            .unwrap();
3057        assert_eq!(
3058            state.read_secure_record(app_tag, key).unwrap().as_deref(),
3059            Some("{\"token\":\"bound\"}")
3060        );
3061        state.delete_secure_record(app_tag, key).unwrap();
3062        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
3063    }
3064
3065    #[test]
3066    fn native_device_proof_fails_closed_without_secure_host_signer() {
3067        let state = OpenRtcTauriState::new(test_config());
3068        assert!(state
3069            .device_public_key("app_native_test")
3070            .expect_err("missing signer must fail")
3071            .contains("secure-storage signer"));
3072    }
3073
3074    fn native_transport_install_context() -> InstallContext {
3075        InstallContext {
3076            data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
3077        }
3078    }
3079
3080    fn managed_session_test_result(device_id: &str) -> StartSessionResult {
3081        StartSessionResult {
3082            local_node_id: format!("node-{device_id}"),
3083            ticket_scope: Some("user-device".to_string()),
3084            ticket: None,
3085            presence_started: true,
3086            auto_connect_started: true,
3087            local_device: openrtc::native_device::NativeDeviceIdentity {
3088                device_id: device_id.to_string(),
3089                device_name: "Test Device".to_string(),
3090                created_at_ms: 1,
3091                updated_at_ms: 1,
3092                name_source: None,
3093                system_info: None,
3094            },
3095        }
3096    }
3097
3098    async fn simulate_managed_session_start(
3099        state: Arc<OpenRtcTauriState>,
3100        device_id: String,
3101        pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
3102        starts: Arc<AtomicUsize>,
3103        stopped_owners: Arc<Mutex<Vec<usize>>>,
3104    ) -> SessionDisposition {
3105        let _start_guard = state.managed_session_start_guard.lock().await;
3106        let client = state.client();
3107        let key = format!("{}:{device_id}", client.app_tag());
3108        let (disposition, active_epoch) = {
3109            let active = state.managed_session.lock().await;
3110            let disposition = managed_session_disposition(
3111                active.as_ref().map(|record| record.key.as_str()),
3112                active
3113                    .as_ref()
3114                    .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
3115                active
3116                    .as_ref()
3117                    .is_some_and(|record| record.result.presence_started),
3118                active
3119                    .as_ref()
3120                    .is_some_and(|record| record.result.auto_connect_started),
3121                &key,
3122                true,
3123                true,
3124            );
3125            (
3126                disposition,
3127                active.as_ref().map(|record| record.owner_epoch),
3128            )
3129        };
3130
3131        if disposition == SessionDisposition::Reuse {
3132            return disposition;
3133        }
3134
3135        if disposition == SessionDisposition::Replace {
3136            let previous = state
3137                .managed_session
3138                .lock()
3139                .await
3140                .take()
3141                .expect("replacement must have a previous owner");
3142            stopped_owners
3143                .lock()
3144                .await
3145                .push(Arc::as_ptr(&previous.owner_client) as usize);
3146        }
3147
3148        starts.fetch_add(1, Ordering::SeqCst);
3149        if let Some((entered, release)) = pause {
3150            entered.send(()).expect("start observer must be waiting");
3151            release.await.expect("start release must be sent");
3152        }
3153
3154        let owner_epoch = match disposition {
3155            SessionDisposition::Refresh => {
3156                active_epoch.expect("refresh must preserve the active owner epoch")
3157            }
3158            SessionDisposition::Start | SessionDisposition::Replace => {
3159                state.allocate_managed_session_owner_epoch()
3160            }
3161            SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
3162        };
3163        let result = managed_session_test_result(&device_id);
3164        *state.managed_session.lock().await = Some(ManagedSessionRecord {
3165            key,
3166            owner_client: client,
3167            owner_epoch,
3168            result,
3169        });
3170        disposition
3171    }
3172
3173    #[tokio::test]
3174    async fn concurrent_same_key_start_reuses_one_owner_epoch() {
3175        let state = Arc::new(OpenRtcTauriState::new(test_config()));
3176        let starts = Arc::new(AtomicUsize::new(0));
3177        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
3178        let (first_entered_tx, first_entered_rx) = oneshot::channel();
3179        let (first_release_tx, first_release_rx) = oneshot::channel();
3180
3181        let first = tokio::spawn(simulate_managed_session_start(
3182            state.clone(),
3183            "same-device".to_string(),
3184            Some((first_entered_tx, first_release_rx)),
3185            starts.clone(),
3186            stopped_owners.clone(),
3187        ));
3188        first_entered_rx.await.expect("first start must pause");
3189
3190        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
3191        let second_state = state.clone();
3192        let second_starts = starts.clone();
3193        let second_stopped_owners = stopped_owners.clone();
3194        let second = tokio::spawn(async move {
3195            second_attempting_tx.send(()).unwrap();
3196            simulate_managed_session_start(
3197                second_state,
3198                "same-device".to_string(),
3199                None,
3200                second_starts,
3201                second_stopped_owners,
3202            )
3203            .await
3204        });
3205        second_attempting_rx.await.unwrap();
3206        tokio::task::yield_now().await;
3207        assert_eq!(starts.load(Ordering::SeqCst), 1);
3208
3209        first_release_tx.send(()).unwrap();
3210        assert_eq!(
3211            first.await.expect("first start task"),
3212            SessionDisposition::Start
3213        );
3214        assert_eq!(
3215            second.await.expect("second start task"),
3216            SessionDisposition::Reuse
3217        );
3218        assert_eq!(starts.load(Ordering::SeqCst), 1);
3219
3220        let active = state
3221            .managed_session
3222            .lock()
3223            .await
3224            .clone()
3225            .expect("active owner");
3226        assert_eq!(active.owner_epoch, 1);
3227        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
3228        assert!(stopped_owners.lock().await.is_empty());
3229    }
3230
3231    #[tokio::test]
3232    async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
3233        let state = Arc::new(OpenRtcTauriState::new(test_config()));
3234        let starts = Arc::new(AtomicUsize::new(0));
3235        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
3236        let (first_entered_tx, first_entered_rx) = oneshot::channel();
3237        let (first_release_tx, first_release_rx) = oneshot::channel();
3238
3239        let first = tokio::spawn(simulate_managed_session_start(
3240            state.clone(),
3241            "first-device".to_string(),
3242            Some((first_entered_tx, first_release_rx)),
3243            starts.clone(),
3244            stopped_owners.clone(),
3245        ));
3246        first_entered_rx.await.expect("first start must pause");
3247        let first_owner = state.client();
3248
3249        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
3250        let second_state = state.clone();
3251        let second_starts = starts.clone();
3252        let second_stopped_owners = stopped_owners.clone();
3253        let second = tokio::spawn(async move {
3254            second_attempting_tx.send(()).unwrap();
3255            simulate_managed_session_start(
3256                second_state,
3257                "second-device".to_string(),
3258                None,
3259                second_starts,
3260                second_stopped_owners,
3261            )
3262            .await
3263        });
3264        second_attempting_rx.await.unwrap();
3265        tokio::task::yield_now().await;
3266        assert!(Arc::ptr_eq(&first_owner, &state.client()));
3267
3268        first_release_tx.send(()).unwrap();
3269        assert_eq!(
3270            first.await.expect("first start task"),
3271            SessionDisposition::Start
3272        );
3273        assert_eq!(
3274            second.await.expect("replacement start task"),
3275            SessionDisposition::Replace
3276        );
3277        assert_eq!(starts.load(Ordering::SeqCst), 2);
3278
3279        let active = state
3280            .managed_session
3281            .lock()
3282            .await
3283            .clone()
3284            .expect("active owner");
3285        assert_eq!(active.owner_epoch, 2);
3286        assert!(active.key.contains("second-device"));
3287        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
3288        assert_eq!(
3289            stopped_owners.lock().await.as_slice(),
3290            [Arc::as_ptr(&first_owner) as usize]
3291        );
3292        assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
3293    }
3294
3295    #[test]
3296    fn derives_app_tag_from_api_key_by_default() {
3297        let config = OpenRtcTauriConfig {
3298            api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
3299            ..OpenRtcTauriConfig::default()
3300        };
3301
3302        assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
3303    }
3304
3305    #[tokio::test]
3306    async fn requested_native_transport_without_installer_keeps_base_route_available() {
3307        let state = OpenRtcTauriState::new(test_config());
3308        let client = state.client();
3309        ensure_requested_native_transports(
3310            &state,
3311            &client,
3312            Some(&requested_ble_config()),
3313            &native_transport_install_context(),
3314        )
3315        .await
3316        .expect("an unavailable optional transport must not fail base Iroh startup");
3317
3318        assert!(state.installed_native_transports.lock().await.is_empty());
3319    }
3320
3321    #[tokio::test]
3322    async fn native_transport_installer_is_idempotent_for_one_client() {
3323        let installs = Arc::new(AtomicUsize::new(0));
3324        let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
3325            Arc::new(FakeTransportInstaller {
3326                installs: installs.clone(),
3327            }),
3328        );
3329        let client = state.client();
3330        let config = requested_ble_config();
3331
3332        ensure_requested_native_transports(
3333            &state,
3334            &client,
3335            Some(&config),
3336            &native_transport_install_context(),
3337        )
3338        .await
3339        .expect("first install");
3340        ensure_requested_native_transports(
3341            &state,
3342            &client,
3343            Some(&config),
3344            &native_transport_install_context(),
3345        )
3346        .await
3347        .expect("idempotent install");
3348
3349        assert_eq!(installs.load(Ordering::SeqCst), 1);
3350    }
3351
3352    #[tokio::test]
3353    async fn unavailable_native_transport_keeps_base_route_available() {
3354        let state = OpenRtcTauriState::new(test_config())
3355            .with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
3356        let client = state.client();
3357
3358        ensure_requested_native_transports(
3359            &state,
3360            &client,
3361            Some(&requested_ble_config()),
3362            &native_transport_install_context(),
3363        )
3364        .await
3365        .expect("optional transport installation failure must not fail base Iroh startup");
3366
3367        assert!(state.installed_native_transports.lock().await.is_empty());
3368    }
3369
3370    #[test]
3371    fn invalid_or_missing_api_key_fails_closed() {
3372        assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
3373        let mut config = OpenRtcTauriConfig::default();
3374        config.api_key = "pk_test_short".to_string();
3375        assert!(config.validated_api_key().is_err());
3376    }
3377
3378    #[test]
3379    fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
3380        let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
3381        let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
3382        let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
3383        let api_key = format!("pk_test_{}", "c".repeat(40));
3384        std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
3385        std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
3386        std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
3387
3388        let config = OpenRtcTauriConfig::from_env();
3389
3390        assert_eq!(config.api_key, api_key);
3391        assert_eq!(
3392            config.app_tag().unwrap(),
3393            openrtc::app_tag_from_api_key(&api_key)
3394        );
3395
3396        match previous_api_key {
3397            Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
3398            None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
3399        }
3400        match previous_project {
3401            Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
3402            None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
3403        }
3404        match previous_app_tag {
3405            Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
3406            None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
3407        }
3408    }
3409
3410    #[test]
3411    fn native_state_has_one_constructor_owned_app_identity() {
3412        let config = test_config();
3413        let expected = openrtc::app_tag_from_api_key(&config.api_key);
3414        let state = OpenRtcTauriState::new(config);
3415        let first = state.client();
3416        let second = state.client();
3417
3418        assert_eq!(first.app_tag(), expected);
3419        assert!(Arc::ptr_eq(&first, &second));
3420    }
3421
3422    #[test]
3423    fn managed_session_uses_persisted_native_device_id() {
3424        let identity = openrtc::native_device::NativeDeviceIdentity {
3425            device_id: "persisted-native-device".to_string(),
3426            device_name: "Mac".to_string(),
3427            created_at_ms: 1,
3428            updated_at_ms: 1,
3429            name_source: None,
3430            system_info: None,
3431        };
3432
3433        assert_eq!(
3434            managed_session_device_id(Some(" desktop-e2e-native "), &identity),
3435            "persisted-native-device"
3436        );
3437        assert_eq!(
3438            managed_session_device_id(None, &identity),
3439            "persisted-native-device"
3440        );
3441    }
3442
3443    #[test]
3444    fn repeated_managed_session_triggers_converge_on_one_owner() {
3445        let key = "app:persisted-native-device";
3446        let cases = [
3447            ("react-remount", true, true, true, true),
3448            ("hmr", true, true, true, true),
3449            ("auth-refresh", true, true, true, true),
3450            ("resume", true, true, true, true),
3451            ("alias-change", true, true, true, true),
3452        ];
3453        for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
3454            assert_eq!(
3455                managed_session_disposition(
3456                    Some(key),
3457                    true,
3458                    active_presence,
3459                    active_auto,
3460                    key,
3461                    requested_presence,
3462                    requested_auto,
3463                ),
3464                SessionDisposition::Reuse,
3465                "{label} must reuse the authoritative tuple"
3466            );
3467        }
3468
3469        assert_eq!(
3470            managed_session_disposition(Some(key), true, false, true, key, true, true),
3471            SessionDisposition::Refresh,
3472            "failed presence startup must retry idempotently"
3473        );
3474        assert_eq!(
3475            managed_session_disposition(
3476                Some(key),
3477                true,
3478                true,
3479                true,
3480                "app:other-native-device",
3481                true,
3482                true,
3483            ),
3484            SessionDisposition::Replace,
3485            "a physical native device change must replace the previous lifecycle owner"
3486        );
3487        assert_eq!(
3488            managed_session_disposition(Some(key), false, true, true, key, true, true),
3489            SessionDisposition::Replace,
3490            "the same tuple on a different client epoch must replace the previous owner"
3491        );
3492    }
3493
3494    #[test]
3495    fn user_device_revocation_invalidates_the_cached_managed_ticket() {
3496        assert!(revokes_managed_session("user-device"));
3497        assert!(revokes_managed_session("  user-device  "));
3498        assert!(!revokes_managed_session("share:example"));
3499        assert!(!revokes_managed_session(""));
3500    }
3501
3502    #[test]
3503    fn managed_session_metadata_uses_authoritative_device_id() {
3504        let raw = serde_json::json!({
3505            "deviceId": "persisted-native-device",
3506            "assistantDevice": {
3507                "deviceId": "desktop-e2e-native"
3508            }
3509        })
3510        .to_string();
3511
3512        let metadata =
3513            metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
3514        let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
3515
3516        assert_eq!(parsed["deviceId"], "desktop-e2e-native");
3517        assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
3518    }
3519}