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 ed25519_dalek::{Signature, VerifyingKey};
8use openrtc::application_crypto_streams::{PeerRecvStream, PeerSendStream};
9use serde::{Deserialize, Serialize};
10#[cfg(feature = "managed-group-encryption")]
11use sha2::{Digest, Sha256};
12use tauri::ipc::{Channel, InvokeBody, Request, Response};
13use tauri::{Manager, Runtime};
14use tokio::sync::Mutex;
15
16#[cfg(feature = "native-broadcast-moq")]
17mod native_broadcast_moq;
18
19const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
20const PEER_BI_STREAM_ID_HEADER: &str = "x-openrtc-stream-id";
21const PEER_ID_HEADER: &str = "x-openrtc-peer-id";
22const FANOUT_CAPABILITY_HEADER: &str = "x-openrtc-fanout-capability";
23const FANOUT_SOURCE_PEER_HEADER: &str = "x-openrtc-fanout-source-peer";
24const FANOUT_APP_TAG_HEADER: &str = "x-openrtc-fanout-app-tag";
25const MEDIA_PUBLICATION_ID_HEADER: &str = "x-openrtc-publication-id";
26const MEDIA_TIMESTAMP_US_HEADER: &str = "x-openrtc-timestamp-us";
27const MEDIA_DURATION_US_HEADER: &str = "x-openrtc-duration-us";
28const MEDIA_KEYFRAME_HEADER: &str = "x-openrtc-keyframe";
29const MEDIA_DISCARDABLE_HEADER: &str = "x-openrtc-discardable";
30
31fn request_body_bytes(request: &Request<'_>) -> Result<Vec<u8>, String> {
32    match request.body() {
33        InvokeBody::Raw(bytes) => Ok(bytes.clone()),
34        // Tauri Android can forward typed-array command bodies through JSON.
35        // Preserve one command contract without falling back to global events.
36        InvokeBody::Json(json) => serde_json::from_value::<Vec<u8>>(json.clone())
37            .map_err(|error| format!("invalid binary IPC payload: {error}")),
38    }
39}
40
41fn required_request_header(request: &Request<'_>, name: &str) -> Result<String, String> {
42    request
43        .headers()
44        .get(name)
45        .and_then(|value| value.to_str().ok())
46        .map(str::trim)
47        .filter(|value| !value.is_empty())
48        .map(ToOwned::to_owned)
49        .ok_or_else(|| format!("{name} header is required"))
50}
51
52fn parsed_request_header<T>(request: &Request<'_>, name: &str) -> Result<T, String>
53where
54    T: std::str::FromStr,
55    T::Err: std::fmt::Display,
56{
57    required_request_header(request, name)?
58        .parse::<T>()
59        .map_err(|error| format!("invalid {name} header: {error}"))
60}
61
62pub type InstallFuture = std::pin::Pin<
63    Box<
64        dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
65            + Send,
66    >,
67>;
68
69pub type BroadcastAdapterFuture<'a> = std::pin::Pin<
70    Box<
71        dyn std::future::Future<
72                Output = Result<Option<openrtc::broadcast::BroadcastAdapterObservation>, String>,
73            > + Send
74            + 'a,
75    >,
76>;
77
78pub type BroadcastObjectsFuture<'a> =
79    std::pin::Pin<Box<dyn std::future::Future<Output = Result<Vec<Vec<u8>>, String>> + Send + 'a>>;
80
81/// Private native relay mechanics boundary. Implementations resolve the
82/// non-secret allocation label to provider credentials inside Rust and never
83/// return those credentials through Tauri IPC.
84pub trait NativeBroadcastAdapter: Send + Sync {
85    fn apply(
86        &self,
87        session: openrtc::broadcast::BroadcastSession,
88        action: openrtc::broadcast::BroadcastAdapterAction,
89    ) -> BroadcastAdapterFuture<'_>;
90
91    /// Drains provider objects for a subscriber. The plugin passes every
92    /// object back through the shared Rust integrity/generation/media decoder
93    /// before returning a high-level media sample over IPC.
94    fn take_inbound_objects(
95        &self,
96        _session: openrtc::broadcast::BroadcastSession,
97        _max: usize,
98    ) -> BroadcastObjectsFuture<'_> {
99        Box::pin(async { Ok(Vec::new()) })
100    }
101}
102
103#[derive(Debug, Clone)]
104pub struct InstallContext {
105    /// Stable host-owned storage used by every native endpoint initializer.
106    /// Installers must derive their Iroh identity from this directory rather
107    /// than creating a transport-specific identity.
108    pub data_dir: PathBuf,
109}
110
111/// Host-supplied implementation for a native transport that cannot be bundled
112/// in the publishable Tauri plugin. The OpenRTC constructor remains the sole
113/// activation switch; registering an installer only declares compiled support.
114pub trait TransportInstaller: Send + Sync {
115    fn id(&self) -> &'static str;
116    fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
117    fn install(
118        &self,
119        client: Arc<openrtc::client::Client>,
120        config: openrtc::client::TransportConfig,
121        context: InstallContext,
122    ) -> InstallFuture;
123}
124
125/// Host-owned private-storage bridge for the OpenRTC 2.0 per-install device
126/// proof key. OpenRTC never receives private key bytes and does not prescribe
127/// an interactive OS credential store. The host exposes only public JWK export
128/// and signing; Plutonium's current shipping owner is its prompt-free,
129/// app-private store.
130pub trait DeviceKeySigner: Send + Sync {
131    fn public_jwk(&self, app_tag: &str) -> Result<serde_json::Value, String>;
132    fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String>;
133
134    /// Host-certified strength of the signer's private-key custody.
135    /// Software is the safe default; only the native host may override it.
136    fn offline_assurance(&self, _app_tag: &str) -> openrtc::offline::OfflineAssurance {
137        openrtc::offline::OfflineAssurance::Software
138    }
139    fn delete(&self, app_tag: &str) -> Result<(), String>;
140    fn read_secure_record(&self, _app_tag: &str, _key: &str) -> Result<Option<String>, String> {
141        Ok(None)
142    }
143    fn write_secure_record(&self, _app_tag: &str, _key: &str, _value: &str) -> Result<(), String> {
144        Err("OpenRTC 2.0 native certificate persistence requires a host secure store".to_string())
145    }
146    fn delete_secure_record(&self, _app_tag: &str, _key: &str) -> Result<(), String> {
147        Ok(())
148    }
149}
150
151struct OfflineDeviceSigner<'a> {
152    signer: &'a dyn DeviceKeySigner,
153    app_tag: &'a str,
154}
155
156impl openrtc::offline::OfflineSigner for OfflineDeviceSigner<'_> {
157    fn verifying_key(&self) -> anyhow::Result<VerifyingKey> {
158        let jwk = self
159            .signer
160            .public_jwk(self.app_tag)
161            .map_err(anyhow::Error::msg)?;
162        validate_public_device_jwk(&jwk).map_err(anyhow::Error::msg)?;
163        let encoded = jwk
164            .get("x")
165            .and_then(serde_json::Value::as_str)
166            .ok_or_else(|| anyhow::anyhow!("OpenRTC device JWK is missing x"))?;
167        let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
168            .decode(encoded)
169            .map_err(|error| anyhow::anyhow!("decode OpenRTC device JWK: {error}"))?;
170        let bytes: [u8; 32] = bytes
171            .try_into()
172            .map_err(|_| anyhow::anyhow!("OpenRTC device JWK must contain 32 key bytes"))?;
173        VerifyingKey::from_bytes(&bytes)
174            .map_err(|error| anyhow::anyhow!("parse OpenRTC device JWK: {error}"))
175    }
176
177    fn sign(&self, message: &[u8]) -> anyhow::Result<Signature> {
178        let bytes = self
179            .signer
180            .sign(self.app_tag, message)
181            .map_err(anyhow::Error::msg)?;
182        Signature::from_slice(&bytes)
183            .map_err(|error| anyhow::anyhow!("parse OpenRTC device signature: {error}"))
184    }
185}
186
187#[derive(Debug, Clone, Serialize)]
188#[serde(rename_all = "camelCase")]
189struct OfflineRuntimeSupport {
190    provisioning: bool,
191    local_mesh: bool,
192    cloud_required: bool,
193    #[serde(skip_serializing_if = "Option::is_none")]
194    reason: Option<&'static str>,
195}
196
197struct InstalledNativeTransport {
198    client_ptr: usize,
199    _runtime: Box<dyn std::any::Any + Send + Sync>,
200}
201
202#[derive(Debug, Clone, Serialize, Deserialize)]
203#[serde(rename_all = "camelCase")]
204pub struct OpenRtcTauriConfig {
205    /// Public OpenRTC developer API key. The app identity is derived locally;
206    /// consumers cannot retarget the native runtime with an app tag or
207    /// provider project identifier.
208    pub api_key: String,
209    #[serde(default)]
210    pub data_dir: Option<PathBuf>,
211    #[serde(default)]
212    pub transport_config: Option<openrtc::client::TransportConfig>,
213}
214
215impl Default for OpenRtcTauriConfig {
216    fn default() -> Self {
217        Self {
218            api_key: String::new(),
219            data_dir: None,
220            transport_config: None,
221        }
222    }
223}
224
225impl OpenRtcTauriConfig {
226    pub fn from_env() -> Self {
227        let api_key = first_env(&[
228            "VITE_OPENRTC_KEY",
229            "VITE_OPENRTC_API_KEY",
230            "VITE_PLUTO_OPENRTC_API_KEY",
231            "OPENRTC_API_KEY",
232        ]);
233        Self {
234            api_key: api_key.unwrap_or_default(),
235            data_dir: None,
236            transport_config: None,
237        }
238    }
239
240    pub fn validated_api_key(&self) -> Result<&str, String> {
241        openrtc::validate_api_key(&self.api_key)
242            .map_err(|error| format!("invalid OpenRTC 2.0 public API key: {error}"))
243    }
244
245    pub fn app_tag(&self) -> Result<String, String> {
246        self.validated_api_key().map(openrtc::app_tag_from_api_key)
247    }
248}
249
250fn first_env(names: &[&str]) -> Option<String> {
251    names.iter().find_map(|name| {
252        std::env::var(name)
253            .ok()
254            .map(|value| value.trim().to_string())
255            .filter(|value| !value.is_empty())
256    })
257}
258
259#[derive(Default)]
260struct TokenRelayState {
261    identity_credential: RwLock<Option<String>>,
262}
263
264impl TokenRelayState {
265    fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
266        let relay = self.clone();
267        Box::new(move || {
268            relay
269                .identity_credential
270                .read()
271                .ok()
272                .and_then(|guard| guard.clone())
273        })
274    }
275
276    fn set(&self, identity_credential: Option<String>) {
277        if let Ok(mut guard) = self.identity_credential.write() {
278            *guard = normalize_token(identity_credential);
279        }
280    }
281}
282
283fn normalize_token(value: Option<String>) -> Option<String> {
284    value
285        .map(|value| value.trim().to_string())
286        .filter(|value| !value.is_empty())
287}
288
289#[derive(Debug, Clone)]
290struct NativeCapabilityRegistration {
291    avenue_kind: String,
292    avenue_id: String,
293    sparse_fanout: bool,
294    requested_architecture: Option<openrtc::native::RoomArchitectureMode>,
295    desired_revision: u64,
296    desired_peers: Vec<serde_json::Value>,
297}
298
299impl NativeCapabilityRegistration {
300    fn uses_sparse_fanout(&self) -> bool {
301        self.requested_architecture
302            .map(|requested| requested == openrtc::native::RoomArchitectureMode::Sparse)
303            .unwrap_or(self.sparse_fanout)
304    }
305}
306
307#[derive(Debug, Default)]
308struct NativeCapabilityRegistry {
309    registrations: HashMap<String, NativeCapabilityRegistration>,
310    root_desired_revision: u64,
311}
312
313impl NativeCapabilityRegistry {
314    #[cfg(test)]
315    fn register(
316        &mut self,
317        capability_key: String,
318        avenue_kind: String,
319        avenue_id: String,
320        sparse_fanout: bool,
321    ) -> Result<(), String> {
322        self.register_with_architecture(capability_key, avenue_kind, avenue_id, sparse_fanout, None)
323    }
324
325    fn register_with_architecture(
326        &mut self,
327        capability_key: String,
328        avenue_kind: String,
329        avenue_id: String,
330        sparse_fanout: bool,
331        requested_architecture: Option<openrtc::native::RoomArchitectureMode>,
332    ) -> Result<(), String> {
333        if requested_architecture.is_some() && avenue_kind != "room" {
334            return Err("native room architecture is valid only for room avenues".to_string());
335        }
336        if let Some(existing) = self.registrations.get(&capability_key) {
337            if existing.avenue_kind == avenue_kind
338                && existing.avenue_id == avenue_id
339                && existing.sparse_fanout == sparse_fanout
340                && existing.requested_architecture == requested_architecture
341            {
342                return Ok(());
343            }
344            return Err(format!(
345                "native capability key {capability_key} is already registered for another avenue"
346            ));
347        }
348        self.registrations.insert(
349            capability_key,
350            NativeCapabilityRegistration {
351                avenue_kind,
352                avenue_id,
353                sparse_fanout,
354                requested_architecture,
355                desired_revision: 0,
356                desired_peers: Vec::new(),
357            },
358        );
359        Ok(())
360    }
361
362    fn capability_keys_for_identity(
363        &self,
364        connection_id: Option<&str>,
365        device_id: Option<&str>,
366        device_id_hint: Option<&str>,
367        remote_node_id: Option<&str>,
368    ) -> Vec<String> {
369        let identities = [connection_id, device_id, device_id_hint, remote_node_id]
370            .into_iter()
371            .flatten()
372            .map(str::trim)
373            .filter(|value| !value.is_empty())
374            .collect::<BTreeSet<_>>();
375        self.registrations
376            .iter()
377            .filter_map(|(key, registration)| {
378                registration
379                    .desired_peers
380                    .iter()
381                    .any(|peer| {
382                        ["connectionId", "deviceId", "nodeId"]
383                            .into_iter()
384                            .filter_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
385                            .map(str::trim)
386                            .any(|value| identities.contains(value))
387                    })
388                    .then(|| key.clone())
389            })
390            .collect::<BTreeSet<_>>()
391            .into_iter()
392            .collect()
393    }
394
395    fn capability_keys_for_state(&self, snapshot: &openrtc::client::StateSnapshot) -> Vec<String> {
396        self.capability_keys_for_identity(
397            Some(&snapshot.connection_id),
398            snapshot.device_id.as_deref(),
399            snapshot.device_id_hint.as_deref(),
400            snapshot.remote_node_id.as_deref(),
401        )
402    }
403
404    fn capability_keys_for_peer_data(
405        &self,
406        event: &openrtc::client::NativePeerDataEvent,
407    ) -> Vec<String> {
408        if let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) {
409            if let Some(capability) = value
410                .get("capability")
411                .and_then(serde_json::Value::as_str)
412                .map(str::trim)
413                .filter(|value| !value.is_empty())
414            {
415                if self.registrations.contains_key(capability) {
416                    return vec![capability.to_string()];
417                }
418                return Vec::new();
419            }
420        }
421        // Peer-data is an application effect, not a connection observation.
422        // It must carry the capability already validated by the Rust ingress
423        // owner; a shared physical peer is never enough to infer one or more
424        // avenue recipients.
425        Vec::new()
426    }
427
428    fn projected_stream_capability(
429        &self,
430        channel: &openrtc::stream_metadata::ChannelMetadata,
431    ) -> Option<String> {
432        let explicit = channel
433            .metadata
434            .as_ref()
435            .and_then(|metadata| metadata.get("openrtcCapability"))
436            .and_then(serde_json::Value::as_str)
437            .map(str::trim)
438            .filter(|value| !value.is_empty());
439        match explicit {
440            Some(key) if self.registrations.contains_key(key) => Some(key.to_string()),
441            _ => None,
442        }
443    }
444
445    fn aggregate_desired_peers(&mut self) -> Result<(u64, String), String> {
446        let role_priority = |value: &serde_json::Value| match value
447            .get("topologyRole")
448            .and_then(serde_json::Value::as_str)
449        {
450            // A non-sparse capability is an unconditional desired reference
451            // and must never be retired by a sparse avenue transition.
452            None => 3_u8,
453            Some("active") => 2,
454            Some("backup") => 1,
455            Some(_) => 0,
456        };
457        let route_evidence = |value: &serde_json::Value| {
458            ["ticket", "nodeId"]
459                .into_iter()
460                .filter(|field| {
461                    value
462                        .get(field)
463                        .and_then(serde_json::Value::as_str)
464                        .is_some_and(|part| !part.trim().is_empty())
465                })
466                .count()
467        };
468        let mut references =
469            std::collections::BTreeMap::<String, Vec<(&str, &serde_json::Value)>>::new();
470        let mut registrations = self.registrations.iter().collect::<Vec<_>>();
471        registrations.sort_by(|(left, _), (right, _)| left.cmp(right));
472        for (capability_key, registration) in registrations {
473            for peer in &registration.desired_peers {
474                let identity = ["deviceId", "nodeId", "ticket"]
475                    .into_iter()
476                    .find_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
477                    .map(str::trim)
478                    .filter(|value| !value.is_empty())
479                    .ok_or_else(|| "desired peer has no stable identity".to_string())?;
480                references
481                    .entry(identity.to_string())
482                    .or_default()
483                    .push((capability_key.as_str(), peer));
484            }
485        }
486        self.root_desired_revision = self.root_desired_revision.saturating_add(1);
487        let peers = references
488            .into_values()
489            .map(|refs| {
490                let semantic_priority = refs
491                    .iter()
492                    .map(|(_, peer)| role_priority(peer))
493                    .max()
494                    .unwrap_or_default();
495                // Pick one complete route atomically. Capability order is the
496                // deterministic tie-breaker; avenue-local records remain
497                // untouched and continue to own their own role/revision.
498                let mut evidence = refs[0];
499                for candidate in refs.iter().copied().skip(1) {
500                    if (route_evidence(candidate.1), role_priority(candidate.1))
501                        > (route_evidence(evidence.1), role_priority(evidence.1))
502                    {
503                        evidence = candidate;
504                    }
505                }
506                let mut peer = evidence.1.clone();
507                if semantic_priority == 3 {
508                    if let Some(object) = peer.as_object_mut() {
509                        object.remove("topologyRole");
510                        object.remove("topologyRevision");
511                    }
512                } else {
513                    peer["topologyRole"] = serde_json::Value::String(
514                        if semantic_priority == 2 {
515                            "active"
516                        } else {
517                            "backup"
518                        }
519                        .to_string(),
520                    );
521                    peer["topologyRevision"] = serde_json::Value::from(self.root_desired_revision);
522                }
523                peer
524            })
525            .collect::<Vec<_>>();
526        let payload = serde_json::to_string(&peers)
527            .map_err(|error| format!("serialize aggregated desired peers: {error}"))?;
528        Ok((self.root_desired_revision, payload))
529    }
530}
531
532fn required_native_capability_part(value: String, label: &str) -> Result<String, String> {
533    let value = value.trim().to_string();
534    if value.is_empty()
535        || value.len() > 192
536        || !value
537            .chars()
538            .all(|character| character.is_ascii_alphanumeric() || "_.:@-".contains(character))
539    {
540        return Err(format!("native {label} is invalid"));
541    }
542    Ok(value)
543}
544
545pub struct OpenRtcTauriState {
546    client: RwLock<Arc<openrtc::client::Client>>,
547    config: OpenRtcTauriConfig,
548    token_relay: Arc<TokenRelayState>,
549    data_dir: Option<PathBuf>,
550    connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
551    peer_data_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
552    presence_loop_active: Mutex<bool>,
553    subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
554    native_projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
555    projected_stream_offered: AtomicU64,
556    projected_stream_decoded: AtomicU64,
557    projected_stream_projected: AtomicU64,
558    projected_stream_unhandled_non_channel: AtomicU64,
559    projected_stream_unhandled_other_channel: AtomicU64,
560    projected_stream_unhandled_no_subscription: AtomicU64,
561    projected_stream_failures: AtomicU64,
562    projected_stream_last_transport_stable_id: AtomicU64,
563    projected_stream_validation_checks: AtomicU64,
564    projected_stream_validation_rejections: AtomicU64,
565    projected_stream_last_validated_transport_stable_id: AtomicU64,
566    projected_stream_authorized: AtomicU64,
567    projected_stream_unauthorized: AtomicU64,
568    projected_stream_last_channel: Mutex<Option<String>>,
569    projected_stream_last_protocol: Mutex<Option<String>>,
570    projected_stream_last_connection_id: Mutex<Option<String>>,
571    projected_stream_last_remote_node_id: Mutex<Option<String>>,
572    peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
573    peer_uni_streams: Mutex<HashMap<String, Arc<Mutex<Option<PeerSendStream>>>>>,
574    portable_media: Mutex<openrtc::media::PortableMediaSession>,
575    broadcast_sessions: Mutex<HashMap<String, openrtc::broadcast::BroadcastSession>>,
576    broadcast_signers: Mutex<HashMap<String, openrtc::broadcast::BroadcastPublisherSigner>>,
577    broadcast_drivers: Mutex<HashMap<String, Arc<Mutex<()>>>>,
578    managed_session_start_guard: Mutex<()>,
579    managed_session_next_owner_epoch: AtomicU64,
580    managed_session: Mutex<Option<ManagedSessionRecord>>,
581    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
582    native_transport_installers: Vec<Arc<dyn TransportInstaller>>,
583    installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
584    native_device_key_signer: Option<Arc<dyn DeviceKeySigner>>,
585    native_broadcast_adapter: Option<Arc<dyn NativeBroadcastAdapter>>,
586    #[cfg(feature = "managed-group-encryption")]
587    managed_groups: Mutex<HashMap<String, openrtc::native::NativeManagedGroupController>>,
588    #[cfg(feature = "native-broadcast-moq")]
589    native_broadcast_moq_adapter: Option<Arc<native_broadcast_moq::NativeBroadcastMoqAdapter>>,
590}
591
592#[derive(Debug, Clone, serde::Serialize)]
593#[serde(rename_all = "camelCase")]
594struct NativeRuntimeStatusResponse {
595    #[serde(flatten)]
596    status: openrtc::client::RuntimeStatus,
597    product_maturity: openrtc::client::ProductCapabilityMaturity,
598}
599
600impl OpenRtcTauriState {
601    pub fn new(config: OpenRtcTauriConfig) -> Self {
602        openrtc::ensure_rustls();
603
604        let token_relay = Arc::new(TokenRelayState::default());
605        let client = build_client(&config, &token_relay)
606            .expect("OpenRTC Tauri 2.0 requires a valid public API key");
607        let data_dir = config.data_dir.clone();
608        #[cfg(feature = "native-broadcast-moq")]
609        let native_broadcast_moq_adapter =
610            Arc::new(native_broadcast_moq::NativeBroadcastMoqAdapter::new());
611
612        Self {
613            client: RwLock::new(client),
614            config,
615            token_relay,
616            data_dir,
617            connection_state_forwarder: std::sync::Mutex::new(None),
618            peer_data_forwarder: std::sync::Mutex::new(None),
619            presence_loop_active: Mutex::new(false),
620            subscriptions: Mutex::new(HashMap::new()),
621            native_projection_subscription: Arc::new(Mutex::new(None)),
622            projected_stream_offered: AtomicU64::new(0),
623            projected_stream_decoded: AtomicU64::new(0),
624            projected_stream_projected: AtomicU64::new(0),
625            projected_stream_unhandled_non_channel: AtomicU64::new(0),
626            projected_stream_unhandled_other_channel: AtomicU64::new(0),
627            projected_stream_unhandled_no_subscription: AtomicU64::new(0),
628            projected_stream_failures: AtomicU64::new(0),
629            projected_stream_last_transport_stable_id: AtomicU64::new(0),
630            projected_stream_validation_checks: AtomicU64::new(0),
631            projected_stream_validation_rejections: AtomicU64::new(0),
632            projected_stream_last_validated_transport_stable_id: AtomicU64::new(0),
633            projected_stream_authorized: AtomicU64::new(0),
634            projected_stream_unauthorized: AtomicU64::new(0),
635            projected_stream_last_channel: Mutex::new(None),
636            projected_stream_last_protocol: Mutex::new(None),
637            projected_stream_last_connection_id: Mutex::new(None),
638            projected_stream_last_remote_node_id: Mutex::new(None),
639            peer_bi_streams: Mutex::new(HashMap::new()),
640            peer_uni_streams: Mutex::new(HashMap::new()),
641            portable_media: Mutex::new(openrtc::media::PortableMediaSession::default()),
642            broadcast_sessions: Mutex::new(HashMap::new()),
643            broadcast_signers: Mutex::new(HashMap::new()),
644            broadcast_drivers: Mutex::new(HashMap::new()),
645            managed_session_start_guard: Mutex::new(()),
646            managed_session_next_owner_epoch: AtomicU64::new(0),
647            managed_session: Mutex::new(None),
648            capability_registry: Arc::new(Mutex::new(NativeCapabilityRegistry::default())),
649            native_transport_installers: Vec::new(),
650            installed_native_transports: Mutex::new(HashMap::new()),
651            native_device_key_signer: None,
652            #[cfg(feature = "managed-group-encryption")]
653            managed_groups: Mutex::new(HashMap::new()),
654            native_broadcast_adapter: {
655                #[cfg(feature = "native-broadcast-moq")]
656                {
657                    Some(native_broadcast_moq_adapter.clone())
658                }
659                #[cfg(not(feature = "native-broadcast-moq"))]
660                {
661                    None
662                }
663            },
664            #[cfg(feature = "native-broadcast-moq")]
665            native_broadcast_moq_adapter: Some(native_broadcast_moq_adapter),
666        }
667    }
668
669    pub fn with_native_transport_installer(
670        mut self,
671        installer: Arc<dyn TransportInstaller>,
672    ) -> Self {
673        self.native_transport_installers.push(installer);
674        self
675    }
676
677    pub fn with_native_device_key_signer(mut self, signer: Arc<dyn DeviceKeySigner>) -> Self {
678        self.native_device_key_signer = Some(signer);
679        self
680    }
681
682    fn offline_runtime_support(&self) -> OfflineRuntimeSupport {
683        let provisioning = self.native_device_key_signer.is_some();
684        OfflineRuntimeSupport {
685            provisioning,
686            local_mesh: cfg!(feature = "transport-lan"),
687            cloud_required: false,
688            reason: if !provisioning {
689                Some("native host device signer is unavailable")
690            } else {
691                (!cfg!(feature = "transport-lan"))
692                    .then_some("native host was built without transport-lan")
693            },
694        }
695    }
696
697    fn project_public_runtime_status(
698        &self,
699        status: openrtc::client::RuntimeStatus,
700    ) -> NativeRuntimeStatusResponse {
701        use openrtc::client::CapabilityMaturity;
702
703        let mut product_maturity = self.client().product_capability_maturity();
704        product_maturity.offline_edge = if self.native_device_key_signer.is_none() {
705            CapabilityMaturity::Unavailable
706        } else if cfg!(feature = "transport-lan") {
707            CapabilityMaturity::Preview
708        } else {
709            CapabilityMaturity::SupportOnly
710        };
711
712        // The Rust broadcast core is present, but openrtc/native deliberately
713        // has no public BroadcastHost adapter in 2.5. Report the shipping
714        // consumer surface, not the private plugin command inventory.
715        product_maturity.broadcast = CapabilityMaturity::Unavailable;
716        NativeRuntimeStatusResponse {
717            status,
718            product_maturity,
719        }
720    }
721
722    pub fn with_native_broadcast_adapter(
723        mut self,
724        adapter: Arc<dyn NativeBroadcastAdapter>,
725    ) -> Self {
726        self.native_broadcast_adapter = Some(adapter);
727        #[cfg(feature = "native-broadcast-moq")]
728        {
729            self.native_broadcast_moq_adapter = None;
730        }
731        self
732    }
733
734    #[cfg(feature = "native-broadcast-moq")]
735    async fn resolve_native_broadcast_access(
736        &self,
737        grant: &str,
738        publication_verifying_key: Option<&[u8; 32]>,
739        expected_issuer_public_key: &[u8; 32],
740    ) -> Result<native_broadcast_moq::NativeBroadcastRelayAccess, String> {
741        let api_key = self.config.validated_api_key()?;
742        let app_tag = self.client().app_tag().to_string();
743        let nonce = uuid::Uuid::new_v4().simple().to_string();
744        let issued_at = std::time::SystemTime::now()
745            .duration_since(std::time::UNIX_EPOCH)
746            .unwrap_or_default()
747            .as_secs();
748        let (challenge, publication_key) = native_broadcast_moq::broadcast_access_challenge(
749            api_key,
750            grant,
751            publication_verifying_key,
752            &nonce,
753            issued_at,
754        );
755        let device_proof = native_broadcast_moq::NativeBroadcastDeviceProof {
756            public_key_jwk: self.device_public_key(&app_tag)?,
757            signature: self.sign_device_proof(&app_tag, &challenge)?,
758            nonce,
759            issued_at,
760        };
761        native_broadcast_moq::resolve_broadcast_access(
762            openrtc::native::OPENRTC_PRODUCTION_CONTROL_PLANE,
763            api_key,
764            grant,
765            publication_key.as_deref(),
766            device_proof,
767            expected_issuer_public_key,
768        )
769        .await
770    }
771
772    fn device_public_key(&self, app_tag: &str) -> Result<serde_json::Value, String> {
773        let app_tag = required_app_tag(app_tag)?;
774        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
775            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
776        })?;
777        let value = signer.public_jwk(app_tag)?;
778        validate_public_device_jwk(&value)?;
779        Ok(value)
780    }
781
782    fn sign_device_proof(&self, app_tag: &str, challenge: &str) -> Result<String, String> {
783        let app_tag = required_app_tag(app_tag)?;
784        if challenge.is_empty() || challenge.len() > 2_048 {
785            return Err("OpenRTC device challenge is invalid".to_string());
786        }
787        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
788            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
789        })?;
790        let signature = signer.sign(app_tag, challenge.as_bytes())?;
791        if signature.len() != 64 {
792            return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
793        }
794        Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
795    }
796
797    fn sign_device_message(&self, app_tag: &str, message: &[u8]) -> Result<String, String> {
798        let app_tag = required_app_tag(app_tag)?;
799        if message.is_empty() || message.len() > 96 * 1024 {
800            return Err("OpenRTC device message is invalid".to_string());
801        }
802        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
803            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
804        })?;
805        let signature = signer.sign(app_tag, message)?;
806        if signature.len() != 64 {
807            return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
808        }
809        Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
810    }
811
812    fn delete_device_key(&self, app_tag: &str) -> Result<(), String> {
813        let app_tag = required_app_tag(app_tag)?;
814        let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
815            "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
816        })?;
817        signer.delete(app_tag)
818    }
819
820    fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
821        let app_tag = required_app_tag(app_tag)?;
822        let key = required_secure_record_key(key)?;
823        self.native_device_key_signer
824            .as_ref()
825            .ok_or_else(|| {
826                "OpenRTC 2.0 native certificate persistence requires a host secure store"
827                    .to_string()
828            })?
829            .read_secure_record(app_tag, key)
830    }
831
832    fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
833        let app_tag = required_app_tag(app_tag)?;
834        let key = required_secure_record_key(key)?;
835        if value.is_empty() || value.len() > 16 * 1024 {
836            return Err("OpenRTC secure record is invalid".to_string());
837        }
838        self.native_device_key_signer
839            .as_ref()
840            .ok_or_else(|| {
841                "OpenRTC 2.0 native certificate persistence requires a host secure store"
842                    .to_string()
843            })?
844            .write_secure_record(app_tag, key, value)
845    }
846
847    fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
848        let app_tag = required_app_tag(app_tag)?;
849        let key = required_secure_record_key(key)?;
850        self.native_device_key_signer
851            .as_ref()
852            .ok_or_else(|| {
853                "OpenRTC 2.0 native certificate persistence requires a host secure store"
854                    .to_string()
855            })?
856            .delete_secure_record(app_tag, key)
857    }
858
859    pub fn client(&self) -> Arc<openrtc::client::Client> {
860        self.client
861            .read()
862            .map(|guard| guard.clone())
863            .expect("OpenRTC Tauri client state is poisoned")
864    }
865
866    pub fn set_identity_credential(&self, identity_credential: Option<String>) {
867        self.token_relay.set(identity_credential);
868    }
869
870    fn allocate_managed_session_owner_epoch(&self) -> u64 {
871        self.managed_session_next_owner_epoch
872            .fetch_add(1, Ordering::SeqCst)
873            + 1
874    }
875
876    fn replace_connection_state_forwarder(&self, client: Arc<openrtc::client::Client>) {
877        let capability_registry = self.capability_registry.clone();
878        let projection_subscription = self.native_projection_subscription.clone();
879        let next = tauri::async_runtime::spawn(async move {
880            forward_connection_state_events(client, capability_registry, projection_subscription)
881                .await;
882        });
883        if let Ok(mut guard) = self.connection_state_forwarder.lock() {
884            if let Some(previous) = guard.replace(next) {
885                previous.abort();
886            }
887        }
888    }
889
890    fn replace_peer_data_forwarder(&self, client: Arc<openrtc::client::Client>) {
891        let capability_registry = self.capability_registry.clone();
892        let projection_subscription = self.native_projection_subscription.clone();
893        let next = tauri::async_runtime::spawn(async move {
894            forward_peer_data_events(client, capability_registry, projection_subscription).await;
895        });
896        if let Ok(mut guard) = self.peer_data_forwarder.lock() {
897            if let Some(previous) = guard.replace(next) {
898                previous.abort();
899            }
900        }
901    }
902}
903
904fn build_client(
905    config: &OpenRtcTauriConfig,
906    token_relay: &Arc<TokenRelayState>,
907) -> Result<Arc<openrtc::client::Client>, String> {
908    let mut builder = openrtc::client::Client::builder(
909        config.validated_api_key()?.to_string(),
910        token_relay.token_provider(),
911    )
912    .map_err(|error| error.to_string())?;
913    if let Some(transport_config) = config.transport_config.clone() {
914        builder = builder.transport_config(transport_config);
915    }
916    Ok(Arc::new(builder.build()))
917}
918
919fn requested_local_device_id(value: Option<&str>) -> Option<String> {
920    value
921        .map(str::trim)
922        .filter(|value| !value.is_empty())
923        .map(ToOwned::to_owned)
924}
925
926fn managed_session_device_id(
927    _requested: Option<&str>,
928    native_identity: &openrtc::native_device::NativeDeviceIdentity,
929) -> String {
930    // The persisted native identity is the only durable-device authority.
931    // Frontend IDs are aliases/display hints and must never create a second
932    // gateway device projection or presence lease for the same native node.
933    native_identity.device_id.clone()
934}
935
936#[cfg(test)]
937fn metadata_with_authoritative_device_id(
938    metadata: Option<String>,
939    device_id: &str,
940) -> Option<String> {
941    let device_id = device_id.trim();
942    if device_id.is_empty() {
943        return metadata;
944    }
945
946    let Some(raw_metadata) = metadata else {
947        return Some(serde_json::json!({ "deviceId": device_id }).to_string());
948    };
949
950    match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
951        Ok(serde_json::Value::Object(mut map)) => {
952            map.insert(
953                "deviceId".to_string(),
954                serde_json::Value::String(device_id.to_string()),
955            );
956            Some(serde_json::Value::Object(map).to_string())
957        }
958        _ => Some(
959            serde_json::json!({
960                "deviceId": device_id,
961                "metadata": raw_metadata,
962            })
963            .to_string(),
964        ),
965    }
966}
967
968struct PeerBiStreamHandle {
969    send: Arc<Mutex<Option<PeerSendStream>>>,
970    recv: Option<PeerRecvStream>,
971    read_task: Option<tokio::task::JoinHandle<()>>,
972}
973
974struct NativeProjectionSubscription {
975    request_id: String,
976    channel: Channel<NativeProjectionEvent>,
977}
978
979#[derive(Serialize, Clone)]
980#[serde(rename_all = "camelCase")]
981struct NativeProjectionEvent {
982    request_id: String,
983    kind: &'static str,
984    capability_keys: Vec<String>,
985    payload: serde_json::Value,
986}
987
988#[derive(Serialize)]
989#[serde(rename_all = "camelCase")]
990pub struct OpenBiResult {
991    stream_id: String,
992    connection_id: Option<String>,
993    remote_node_id: String,
994}
995
996#[derive(Serialize)]
997#[serde(rename_all = "camelCase")]
998pub struct OpenUniResult {
999    stream_id: String,
1000    connection_id: Option<String>,
1001    remote_node_id: String,
1002}
1003
1004#[derive(Serialize, Clone)]
1005#[serde(rename_all = "camelCase")]
1006struct IncomingPeerBiStreamEvent {
1007    request_id: String,
1008    stream_id: String,
1009    connection_id: Option<String>,
1010    remote_node_id: String,
1011    transport_stable_id: u64,
1012    channel: Option<openrtc::stream_metadata::ChannelMetadata>,
1013    application_authorized: bool,
1014    capability_key: String,
1015}
1016
1017#[derive(Debug, Clone, Serialize)]
1018#[serde(rename_all = "camelCase")]
1019struct ProjectedStreamDiagnostics {
1020    subscription_active: bool,
1021    offered: u64,
1022    decoded: u64,
1023    projected: u64,
1024    unhandled_non_channel: u64,
1025    unhandled_other_channel: u64,
1026    unhandled_no_subscription: u64,
1027    failures: u64,
1028    last_transport_stable_id: u64,
1029    validation_checks: u64,
1030    validation_rejections: u64,
1031    last_validated_transport_stable_id: u64,
1032    authorized: u64,
1033    unauthorized: u64,
1034    last_channel: Option<String>,
1035    last_protocol: Option<String>,
1036    last_connection_id: Option<String>,
1037    last_remote_node_id: Option<String>,
1038}
1039
1040/// Result of offering one already-admitted, already-protected peer stream to
1041/// the OpenRTC-owned native protocol projection. A host must continue its own
1042/// protocol dispatch only for `Unhandled`; there is still exactly one owner of
1043/// the underlying incoming-stream queue.
1044pub enum ProjectedPeerBiHandoff {
1045    Handled,
1046    Unhandled {
1047        send: PeerSendStream,
1048        recv: PeerRecvStream,
1049    },
1050}
1051
1052#[derive(Serialize, Clone)]
1053#[serde(rename_all = "camelCase")]
1054struct PeerBiStreamClosedEvent {
1055    r#type: &'static str,
1056    error: Option<String>,
1057}
1058
1059#[derive(Debug, Clone, Serialize)]
1060#[serde(rename_all = "camelCase")]
1061pub struct StartSessionResult {
1062    local_node_id: String,
1063    ticket_scope: Option<String>,
1064    ticket: Option<String>,
1065    presence_started: bool,
1066    auto_connect_started: bool,
1067    local_device: openrtc::native_device::NativeDeviceIdentity,
1068}
1069
1070#[derive(Clone)]
1071struct ManagedSessionRecord {
1072    key: String,
1073    owner_client: Arc<openrtc::client::Client>,
1074    owner_epoch: u64,
1075    result: StartSessionResult,
1076}
1077
1078#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1079enum SessionDisposition {
1080    Start,
1081    Reuse,
1082    Refresh,
1083    Replace,
1084}
1085
1086fn managed_session_disposition(
1087    active_key: Option<&str>,
1088    active_client_matches: bool,
1089    active_presence_started: bool,
1090    active_auto_connect_started: bool,
1091    requested_key: &str,
1092    requested_presence: bool,
1093    requested_auto_connect: bool,
1094) -> SessionDisposition {
1095    let Some(active_key) = active_key else {
1096        return SessionDisposition::Start;
1097    };
1098    if !active_client_matches || active_key != requested_key {
1099        return SessionDisposition::Replace;
1100    }
1101    if (!requested_presence || active_presence_started)
1102        && (!requested_auto_connect || active_auto_connect_started)
1103    {
1104        SessionDisposition::Reuse
1105    } else {
1106        SessionDisposition::Refresh
1107    }
1108}
1109
1110fn revokes_managed_session(scope: &str) -> bool {
1111    scope.trim() == "user-device"
1112}
1113
1114pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
1115    init_with_state(OpenRtcTauriState::new(config))
1116}
1117
1118pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
1119    tauri::plugin::Builder::new(PLUGIN_NAME)
1120        .setup(move |app, _api| {
1121            let client = state.client();
1122            let transport_config = state.config.transport_config.clone();
1123            // Optional native transports must decorate the endpoint before the
1124            // host app can initialize the shared Iroh client during its setup.
1125            // Managed-session startup is too late for custom transports because
1126            // Iroh endpoint hooks are immutable after bind/adoption.
1127            let install_context = InstallContext {
1128                data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
1129            };
1130            tauri::async_runtime::block_on(ensure_requested_native_transports(
1131                &state,
1132                &client,
1133                transport_config.as_ref(),
1134                &install_context,
1135            ))
1136            .map_err(std::io::Error::other)?;
1137            state.replace_connection_state_forwarder(client.clone());
1138            state.replace_peer_data_forwarder(client);
1139            app.manage(state);
1140            Ok(())
1141        })
1142        .invoke_handler(tauri::generate_handler![
1143            openrtc_set_identity_credential,
1144            openrtc_device_public_key,
1145            openrtc_sign_device_proof,
1146            openrtc_delete_device_key,
1147            openrtc_read_secure_record,
1148            openrtc_write_secure_record,
1149            openrtc_delete_secure_record,
1150            openrtc_offline_runtime_support,
1151            openrtc_create_offline_enrollment_request,
1152            openrtc_verify_offline_enrollment_request,
1153            rtc_native_status,
1154            get_rtc_local_device_info,
1155            update_rtc_local_device_name,
1156            get_iroh_node_id,
1157            start_iroh_node,
1158            get_iroh_endpoint_ticket,
1159            register_session_token,
1160            get_endpoint_ticket_with_token,
1161            validate_session_token,
1162            revoke_session_tokens_by_scope,
1163            stop_rtc_presence_loop,
1164            register_rtc_capability,
1165            openrtc_init_managed_room_group,
1166            openrtc_handle_managed_room_prepare_page,
1167            openrtc_handle_managed_room_artifact_chunk,
1168            openrtc_seal_managed_room_payload,
1169            openrtc_open_managed_room_payload,
1170            openrtc_forget_managed_room_group,
1171            unregister_rtc_capability,
1172            start_rtc_managed_session,
1173            start_rtc_external_auto_connect,
1174            submit_rtc_desired_peers,
1175            stop_rtc_auto_connect,
1176            notify_rtc_network_change,
1177            set_rtc_transport_priority,
1178            connect_to_device,
1179            disconnect_device,
1180            set_auto_connect_excluded,
1181            set_rtc_external_auto_connect_excluded,
1182            resolve_rtc_peer_connection_records,
1183            resolve_rtc_peer_identity,
1184            get_rtc_peer_session,
1185            list_rtc_peer_sessions,
1186            list_rtc_managed_connections,
1187            wait_for_rtc_settled_peer,
1188            list_rtc_connection_states,
1189            get_rtc_connection_state,
1190            stop_rtc_subscription,
1191            start_native_projection,
1192            get_projected_stream_diagnostics,
1193            is_current_transport_stable_id,
1194            send_peer_message,
1195            encode_sparse_fanout_message,
1196            accept_sparse_fanout_message,
1197            sparse_fanout_diagnostics,
1198            prepare_openrtc_broadcast_publisher,
1199            release_openrtc_broadcast_publisher,
1200            open_openrtc_broadcast,
1201            begin_openrtc_broadcast_publication,
1202            publish_openrtc_broadcast_sample,
1203            receive_openrtc_broadcast_media,
1204            revoke_openrtc_broadcast,
1205            close_openrtc_broadcast,
1206            record_sparse_fanout_forward_queue_drop,
1207            is_peer_connected,
1208            open_peer_bi_stream,
1209            open_peer_bi_transport_only_stream,
1210            open_peer_uni_stream,
1211            write_peer_bi_stream,
1212            start_peer_bi_stream_read,
1213            cancel_peer_bi_stream_read,
1214            close_peer_bi_stream,
1215            write_peer_uni_stream,
1216            close_peer_uni_stream,
1217            begin_openrtc_media_publication,
1218            encode_openrtc_media_sample,
1219            pause_openrtc_media_publication,
1220            retire_openrtc_media_publication,
1221            retire_openrtc_media_receiver,
1222            decode_openrtc_media_chunk,
1223            decode_openrtc_media_control,
1224        ])
1225        .build()
1226}
1227
1228pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
1229    init(OpenRtcTauriConfig::from_env())
1230}
1231
1232async fn send_native_projection<T: Serialize>(
1233    subscription: &Arc<Mutex<Option<NativeProjectionSubscription>>>,
1234    kind: &'static str,
1235    capability_keys: Vec<String>,
1236    payload: &T,
1237) {
1238    if capability_keys.is_empty() {
1239        return;
1240    }
1241    let Ok(payload) = serde_json::to_value(payload) else {
1242        return;
1243    };
1244    let target = subscription.lock().await.as_ref().map(|subscription| {
1245        (
1246            subscription.request_id.clone(),
1247            subscription.channel.clone(),
1248        )
1249    });
1250    let Some((request_id, channel)) = target else {
1251        return;
1252    };
1253    let _ = channel.send(NativeProjectionEvent {
1254        request_id,
1255        kind,
1256        capability_keys,
1257        payload,
1258    });
1259}
1260
1261async fn forward_connection_state_events(
1262    client: Arc<openrtc::client::Client>,
1263    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
1264    projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
1265) {
1266    let mut rx = client.connection_state_updates();
1267    while let Ok(snapshot) = rx.recv().await {
1268        let keys = capability_registry
1269            .lock()
1270            .await
1271            .capability_keys_for_state(&snapshot);
1272        send_native_projection(
1273            &projection_subscription,
1274            "connection-state",
1275            keys,
1276            &snapshot,
1277        )
1278        .await;
1279    }
1280}
1281
1282async fn forward_peer_data_events(
1283    client: Arc<openrtc::client::Client>,
1284    capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
1285    projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
1286) {
1287    let mut rx = client.subscribe_native_peer_data();
1288    while let Ok(event) = rx.recv().await {
1289        let keys = capability_registry
1290            .lock()
1291            .await
1292            .capability_keys_for_peer_data(&event);
1293        send_native_projection(&projection_subscription, "peer-data", keys, &event).await;
1294    }
1295}
1296
1297async fn replay_current_connection_states(state: &OpenRtcTauriState) {
1298    let snapshots = state.client().connection_states().await;
1299    let registry = state.capability_registry.lock().await;
1300    for snapshot in snapshots {
1301        let keys = registry.capability_keys_for_state(&snapshot);
1302        send_native_projection(
1303            &state.native_projection_subscription,
1304            "connection-state",
1305            keys,
1306            &snapshot,
1307        )
1308        .await;
1309    }
1310}
1311
1312fn app_data_dir<R: Runtime>(
1313    app: &tauri::AppHandle<R>,
1314    state: &OpenRtcTauriState,
1315) -> Result<PathBuf, String> {
1316    state
1317        .data_dir
1318        .clone()
1319        .or_else(|| app.path().app_data_dir().ok())
1320        .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
1321}
1322
1323async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
1324    if let Some(node_id) = client.current_node_id().await {
1325        return Ok(node_id);
1326    }
1327    client
1328        .init_iroh(None, Vec::new())
1329        .await
1330        .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
1331}
1332
1333async fn ensure_requested_native_transports(
1334    state: &OpenRtcTauriState,
1335    client: &Arc<openrtc::client::Client>,
1336    transports: Option<&openrtc::client::TransportConfig>,
1337    context: &InstallContext,
1338) -> Result<(), String> {
1339    let Some(transports) = transports else {
1340        return Ok(());
1341    };
1342    let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
1343    if ble_requested
1344        && !state
1345            .native_transport_installers
1346            .iter()
1347            .any(|installer| installer.is_requested(transports))
1348    {
1349        eprintln!(
1350            "[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
1351        );
1352        return Ok(());
1353    }
1354    let client_ptr = Arc::as_ptr(client) as usize;
1355    for installer in &state.native_transport_installers {
1356        if !installer.is_requested(transports) {
1357            continue;
1358        }
1359        let already_installed = state
1360            .installed_native_transports
1361            .lock()
1362            .await
1363            .get(installer.id())
1364            .is_some_and(|installed| installed.client_ptr == client_ptr);
1365        if already_installed {
1366            continue;
1367        }
1368        if client.current_node_id().await.is_some() {
1369            eprintln!(
1370                "[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
1371                installer.id()
1372            );
1373            continue;
1374        }
1375        let runtime = match installer
1376            .install(client.clone(), transports.clone(), context.clone())
1377            .await
1378        {
1379            Ok(runtime) => runtime,
1380            Err(error) => {
1381                eprintln!(
1382                    "[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
1383                    installer.id(),
1384                    error
1385                );
1386                continue;
1387            }
1388        };
1389        state.installed_native_transports.lock().await.insert(
1390            installer.id(),
1391            InstalledNativeTransport {
1392                client_ptr,
1393                _runtime: runtime,
1394            },
1395        );
1396    }
1397    Ok(())
1398}
1399
1400async fn register_peer_bi_stream(
1401    state: &OpenRtcTauriState,
1402    connection_id: Option<String>,
1403    remote_node_id: String,
1404    send: PeerSendStream,
1405    recv: PeerRecvStream,
1406) -> OpenBiResult {
1407    let stream_id = uuid::Uuid::new_v4().to_string();
1408    let send = Arc::new(Mutex::new(Some(send)));
1409
1410    state.peer_bi_streams.lock().await.insert(
1411        stream_id.clone(),
1412        PeerBiStreamHandle {
1413            send,
1414            // Streams stay paused until JavaScript supplies its ordered Tauri
1415            // Channel. This prevents byte zero from racing listener setup.
1416            recv: Some(recv),
1417            read_task: None,
1418        },
1419    );
1420
1421    OpenBiResult {
1422        stream_id,
1423        connection_id,
1424        remote_node_id,
1425    }
1426}
1427
1428async fn read_peer_exact(recv: &mut PeerRecvStream, buffer: &mut [u8]) -> std::io::Result<()> {
1429    let mut offset = 0;
1430    while offset < buffer.len() {
1431        let read = recv.read(&mut buffer[offset..]).await?;
1432        if read == 0 {
1433            return Err(std::io::Error::new(
1434                std::io::ErrorKind::UnexpectedEof,
1435                "peer stream ended during channel classification",
1436            ));
1437        }
1438        offset += read;
1439    }
1440    Ok(())
1441}
1442
1443fn is_openrtc_projected_channel(channel: &openrtc::stream_metadata::ChannelMetadata) -> bool {
1444    channel.channel_id == openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID
1445}
1446
1447/// Offer an already-admitted and application-protected stream from the
1448/// embedding native host to OpenRTC's typed IPC projection.
1449///
1450/// This is the integration point for Tauri applications that already consume
1451/// the client's sole `incoming_streams()` queue for product protocols. The
1452/// helper performs bounded channel-envelope classification and restores every
1453/// consumed plaintext byte when the stream is not owned by OpenRTC.
1454pub async fn handoff_projected_peer_bi_stream<R: Runtime>(
1455    app: tauri::AppHandle<R>,
1456    connection_id: Option<String>,
1457    remote_node_id: String,
1458    transport_stable_id: u64,
1459    send: PeerSendStream,
1460    mut recv: PeerRecvStream,
1461) -> Result<ProjectedPeerBiHandoff, String> {
1462    let state = app.state::<OpenRtcTauriState>();
1463    state
1464        .projected_stream_offered
1465        .fetch_add(1, Ordering::Relaxed);
1466    state
1467        .projected_stream_last_transport_stable_id
1468        .store(transport_stable_id, Ordering::Relaxed);
1469    *state.projected_stream_last_connection_id.lock().await = connection_id.clone();
1470    *state.projected_stream_last_remote_node_id.lock().await = Some(remote_node_id.clone());
1471    let mut envelope = Vec::new();
1472    let channel = loop {
1473        match openrtc::stream_metadata::decode_prefix(&envelope) {
1474            openrtc::stream_metadata::DecodeDecision::NotMatched => {
1475                state
1476                    .projected_stream_unhandled_non_channel
1477                    .fetch_add(1, Ordering::Relaxed);
1478                return Ok(ProjectedPeerBiHandoff::Unhandled {
1479                    send,
1480                    recv: recv.with_plaintext_prefix(envelope),
1481                });
1482            }
1483            openrtc::stream_metadata::DecodeDecision::Decoded(prefix) => break prefix.channel,
1484            openrtc::stream_metadata::DecodeDecision::NeedMore(required) => {
1485                let missing = required.saturating_sub(envelope.len());
1486                if missing == 0 {
1487                    return Err("channel-envelope decoder made no progress".to_string());
1488                }
1489                let mut bytes = vec![0u8; missing];
1490                if let Err(error) = read_peer_exact(&mut recv, &mut bytes).await {
1491                    state
1492                        .projected_stream_failures
1493                        .fetch_add(1, Ordering::Relaxed);
1494                    return Err(format!("read projected channel envelope: {error}"));
1495                }
1496                envelope.extend_from_slice(&bytes);
1497            }
1498        }
1499    };
1500    state
1501        .projected_stream_decoded
1502        .fetch_add(1, Ordering::Relaxed);
1503    *state.projected_stream_last_channel.lock().await = Some(channel.channel_id.clone());
1504    *state.projected_stream_last_protocol.lock().await = channel
1505        .metadata
1506        .as_ref()
1507        .and_then(|metadata| metadata.get("protocol"))
1508        .and_then(serde_json::Value::as_str)
1509        .map(str::to_string);
1510
1511    if !is_openrtc_projected_channel(&channel) {
1512        state
1513            .projected_stream_unhandled_other_channel
1514            .fetch_add(1, Ordering::Relaxed);
1515        return Ok(ProjectedPeerBiHandoff::Unhandled {
1516            send,
1517            recv: recv.with_plaintext_prefix(envelope),
1518        });
1519    }
1520
1521    let Some(capability_key) = state
1522        .capability_registry
1523        .lock()
1524        .await
1525        .projected_stream_capability(&channel)
1526    else {
1527        // Capability ownership is explicit on the protected channel envelope.
1528        // Unknown or unscoped descriptors remain available to the embedding
1529        // host; they are never assigned by listener order or capability count.
1530        state
1531            .projected_stream_unhandled_other_channel
1532            .fetch_add(1, Ordering::Relaxed);
1533        return Ok(ProjectedPeerBiHandoff::Unhandled {
1534            send,
1535            recv: recv.with_plaintext_prefix(envelope),
1536        });
1537    };
1538
1539    let subscription = state.native_projection_subscription.lock().await;
1540    let Some((request_id, projection_channel)) = subscription.as_ref().map(|subscription| {
1541        (
1542            subscription.request_id.clone(),
1543            subscription.channel.clone(),
1544        )
1545    }) else {
1546        state
1547            .projected_stream_unhandled_no_subscription
1548            .fetch_add(1, Ordering::Relaxed);
1549        return Ok(ProjectedPeerBiHandoff::Unhandled {
1550            send,
1551            recv: recv.with_plaintext_prefix(envelope),
1552        });
1553    };
1554    drop(subscription);
1555
1556    let application_authorized = connection_id.as_deref().is_some_and(|connection_id| {
1557        matches!(
1558            state.client().session_admission(connection_id),
1559            openrtc::session_token::SessionAdmission::Accepted { .. }
1560        )
1561    });
1562    if application_authorized {
1563        state
1564            .projected_stream_authorized
1565            .fetch_add(1, Ordering::Relaxed);
1566    } else {
1567        state
1568            .projected_stream_unauthorized
1569            .fetch_add(1, Ordering::Relaxed);
1570    }
1571
1572    let result = register_peer_bi_stream(
1573        state.inner(),
1574        connection_id,
1575        remote_node_id.clone(),
1576        send,
1577        recv,
1578    )
1579    .await;
1580    let event = IncomingPeerBiStreamEvent {
1581        request_id: request_id.clone(),
1582        stream_id: result.stream_id.clone(),
1583        connection_id: result.connection_id,
1584        remote_node_id,
1585        transport_stable_id,
1586        channel: Some(channel),
1587        application_authorized,
1588        capability_key: capability_key.clone(),
1589    };
1590    let payload = serde_json::to_value(event)
1591        .map_err(|error| format!("serialize projected peer stream: {error}"))?;
1592    if let Err(error) = projection_channel.send(NativeProjectionEvent {
1593        request_id,
1594        kind: "stream",
1595        capability_keys: vec![capability_key],
1596        payload,
1597    }) {
1598        state.peer_bi_streams.lock().await.remove(&result.stream_id);
1599        state
1600            .projected_stream_failures
1601            .fetch_add(1, Ordering::Relaxed);
1602        return Err(format!("send projected peer stream: {error}"));
1603    }
1604    state
1605        .projected_stream_projected
1606        .fetch_add(1, Ordering::Relaxed);
1607    Ok(ProjectedPeerBiHandoff::Handled)
1608}
1609
1610fn native_stream_trace_enabled() -> bool {
1611    cfg!(debug_assertions)
1612        || std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1")
1613}
1614
1615fn spawn_peer_bi_stream_reader(
1616    channel: Channel<Response>,
1617    stream_id: String,
1618    mut recv: PeerRecvStream,
1619) -> tokio::task::JoinHandle<()> {
1620    tokio::spawn(async move {
1621        let mut chunk = vec![0_u8; 64 * 1024];
1622        let mut close_error: Option<String> = None;
1623        loop {
1624            match recv.read(&mut chunk).await {
1625                Ok(0) => break,
1626                Ok(n) => {
1627                    if native_stream_trace_enabled() {
1628                        eprintln!(
1629                            "[openrtc-tauri][peer-bi-stream] stream_id={} phase=plaintext-chunk bytes={}",
1630                            stream_id, n
1631                        );
1632                    }
1633                    if channel.send(Response::new(chunk[..n].to_vec())).is_err() {
1634                        close_error = Some("native stream IPC channel closed".to_string());
1635                        break;
1636                    }
1637                }
1638                Err(error) => {
1639                    if native_stream_trace_enabled() {
1640                        eprintln!(
1641                            "[openrtc-tauri][peer-bi-stream] stream_id={} phase=read-error error={}",
1642                            stream_id, error
1643                        );
1644                    }
1645                    close_error = Some(error.to_string());
1646                    break;
1647                }
1648            }
1649        }
1650        let close = serde_json::to_string(&PeerBiStreamClosedEvent {
1651            r#type: "closed",
1652            error: close_error,
1653        })
1654        .expect("peer stream close event serializes");
1655        let _ = channel.send(Response::new(close));
1656    })
1657}
1658
1659async fn open_peer_bi_with<R: Runtime, F, Fut>(
1660    _app: tauri::AppHandle<R>,
1661    state: tauri::State<'_, OpenRtcTauriState>,
1662    peer_id: String,
1663    timeout_ms: Option<u64>,
1664    open: F,
1665) -> Result<OpenBiResult, String>
1666where
1667    F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
1668    Fut: std::future::Future<
1669        Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
1670    >,
1671{
1672    let peer_id = peer_id.trim().to_string();
1673    if peer_id.is_empty() {
1674        return Err("peerId is required".to_string());
1675    }
1676
1677    let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
1678        .await
1679        .map_err(|error| format!("open peer bi stream failed: {error}"))?;
1680    Ok(register_peer_bi_stream(&state, connection_id, remote_node_id, send, recv).await)
1681}
1682
1683/// Relay an opaque OpenRTC-issued identity credential into the provider-neutral
1684/// Rust runtime. The host owns acquisition and renewal; this command never
1685/// interprets the value as a Firebase token and accepts no refresh token.
1686#[tauri::command]
1687async fn openrtc_set_identity_credential(
1688    state: tauri::State<'_, OpenRtcTauriState>,
1689    identity_credential: Option<String>,
1690) -> Result<(), String> {
1691    state.set_identity_credential(identity_credential);
1692    if *state.presence_loop_active.lock().await {
1693        state.client().request_presence_update();
1694    }
1695    Ok(())
1696}
1697
1698fn required_app_tag(value: &str) -> Result<&str, String> {
1699    let value = value.trim();
1700    if value.len() < 5
1701        || value.len() > 80
1702        || !value.bytes().all(|byte| {
1703            byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
1704        })
1705    {
1706        return Err("OpenRTC app tag is invalid".to_string());
1707    }
1708    Ok(value)
1709}
1710
1711fn validate_public_device_jwk(value: &serde_json::Value) -> Result<(), String> {
1712    let object = value
1713        .as_object()
1714        .ok_or_else(|| "OpenRTC device public key is invalid".to_string())?;
1715    let x = object
1716        .get("x")
1717        .and_then(serde_json::Value::as_str)
1718        .unwrap_or_default();
1719    if object.get("kty").and_then(serde_json::Value::as_str) != Some("OKP")
1720        || object.get("crv").and_then(serde_json::Value::as_str) != Some("Ed25519")
1721        || object.contains_key("d")
1722        || x.len() != 43
1723        || !x
1724            .bytes()
1725            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
1726    {
1727        return Err("OpenRTC device public key must be a public Ed25519 JWK".to_string());
1728    }
1729    Ok(())
1730}
1731
1732fn required_secure_record_key(value: &str) -> Result<&str, String> {
1733    let value = value.trim();
1734    if value.is_empty()
1735        || value.len() > 512
1736        || !value
1737            .bytes()
1738            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':'))
1739    {
1740        return Err("OpenRTC secure record key is invalid".to_string());
1741    }
1742    Ok(value)
1743}
1744
1745#[tauri::command]
1746async fn openrtc_device_public_key(
1747    state: tauri::State<'_, OpenRtcTauriState>,
1748    app_tag: String,
1749) -> Result<serde_json::Value, String> {
1750    state.device_public_key(&app_tag)
1751}
1752
1753#[tauri::command]
1754async fn openrtc_sign_device_proof(
1755    state: tauri::State<'_, OpenRtcTauriState>,
1756    app_tag: String,
1757    challenge: String,
1758) -> Result<String, String> {
1759    state.sign_device_proof(&app_tag, &challenge)
1760}
1761
1762#[tauri::command]
1763async fn openrtc_delete_device_key(
1764    state: tauri::State<'_, OpenRtcTauriState>,
1765    app_tag: String,
1766) -> Result<(), String> {
1767    state.delete_device_key(&app_tag)
1768}
1769
1770#[tauri::command]
1771async fn openrtc_read_secure_record(
1772    state: tauri::State<'_, OpenRtcTauriState>,
1773    app_tag: String,
1774    key: String,
1775) -> Result<Option<String>, String> {
1776    state.read_secure_record(&app_tag, &key)
1777}
1778
1779#[tauri::command]
1780async fn openrtc_write_secure_record(
1781    state: tauri::State<'_, OpenRtcTauriState>,
1782    app_tag: String,
1783    key: String,
1784    value: String,
1785) -> Result<(), String> {
1786    state.write_secure_record(&app_tag, &key, &value)
1787}
1788
1789#[tauri::command]
1790async fn openrtc_delete_secure_record(
1791    state: tauri::State<'_, OpenRtcTauriState>,
1792    app_tag: String,
1793    key: String,
1794) -> Result<(), String> {
1795    state.delete_secure_record(&app_tag, &key)
1796}
1797
1798#[tauri::command]
1799async fn openrtc_offline_runtime_support(
1800    state: tauri::State<'_, OpenRtcTauriState>,
1801) -> Result<OfflineRuntimeSupport, String> {
1802    Ok(state.offline_runtime_support())
1803}
1804
1805#[tauri::command]
1806async fn openrtc_create_offline_enrollment_request(
1807    state: tauri::State<'_, OpenRtcTauriState>,
1808    trust_domain: String,
1809    device_id: String,
1810    endpoint_id: String,
1811    enrollment_nonce: String,
1812    requested_roles: Vec<String>,
1813) -> Result<openrtc::offline::OfflineEnrollmentRequest, String> {
1814    let client = state.client();
1815    let current_endpoint_id = client.current_node_id().await.ok_or_else(|| {
1816        "OpenRTC Iroh endpoint must be started before offline enrollment".to_string()
1817    })?;
1818    if current_endpoint_id != endpoint_id.trim() {
1819        return Err(
1820            "offline enrollment endpoint does not match the native Rust endpoint".to_string(),
1821        );
1822    }
1823    let app_tag = client.app_tag();
1824    let signer = state
1825        .native_device_key_signer
1826        .as_deref()
1827        .ok_or_else(|| "offline enrollment requires the native host device signer".to_string())?;
1828    openrtc::offline::OfflineEnrollmentRequest::create(
1829        &OfflineDeviceSigner { signer, app_tag },
1830        &trust_domain,
1831        &device_id,
1832        &current_endpoint_id,
1833        &enrollment_nonce,
1834        requested_roles,
1835        signer.offline_assurance(app_tag),
1836        openrtc::session_token::now_unix_ms(),
1837    )
1838    .map_err(|error| error.to_string())
1839}
1840
1841#[tauri::command]
1842async fn openrtc_verify_offline_enrollment_request(
1843    request: openrtc::offline::OfflineEnrollmentRequest,
1844) -> Result<(), String> {
1845    request
1846        .verify()
1847        .map(|_| ())
1848        .map_err(|error| error.to_string())
1849}
1850
1851#[tauri::command]
1852async fn rtc_native_status(
1853    state: tauri::State<'_, OpenRtcTauriState>,
1854) -> Result<NativeRuntimeStatusResponse, String> {
1855    Ok(state.project_public_runtime_status(state.client().runtime_status().await))
1856}
1857
1858#[tauri::command]
1859async fn get_rtc_local_device_info<R: Runtime>(
1860    app: tauri::AppHandle<R>,
1861    state: tauri::State<'_, OpenRtcTauriState>,
1862) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1863    let data_dir = app_data_dir(&app, &state)?;
1864    state
1865        .client()
1866        .init_native_device_identity(data_dir, None)
1867        .await
1868        .map_err(|error| error.to_string())
1869}
1870
1871#[tauri::command]
1872async fn update_rtc_local_device_name<R: Runtime>(
1873    app: tauri::AppHandle<R>,
1874    state: tauri::State<'_, OpenRtcTauriState>,
1875    device_name: String,
1876) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1877    let data_dir = app_data_dir(&app, &state)?;
1878    let _ = state
1879        .client()
1880        .init_native_device_identity(data_dir, None)
1881        .await
1882        .map_err(|error| error.to_string())?;
1883    state
1884        .client()
1885        .update_native_device_name(&device_name)
1886        .await
1887        .map_err(|error| error.to_string())
1888}
1889
1890#[tauri::command]
1891async fn get_iroh_node_id(
1892    state: tauri::State<'_, OpenRtcTauriState>,
1893) -> Result<Option<String>, String> {
1894    Ok(state.client().current_node_id().await)
1895}
1896
1897#[tauri::command]
1898async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
1899    ensure_iroh_node(state.client()).await
1900}
1901
1902#[tauri::command]
1903async fn get_iroh_endpoint_ticket(
1904    state: tauri::State<'_, OpenRtcTauriState>,
1905) -> Result<String, String> {
1906    ensure_iroh_node(state.client()).await?;
1907    state
1908        .client()
1909        .endpoint_ticket()
1910        .await
1911        .map_err(|error| error.to_string())
1912}
1913
1914#[tauri::command]
1915async fn register_session_token(
1916    state: tauri::State<'_, OpenRtcTauriState>,
1917    token: String,
1918    scope: String,
1919    max_connections: u32,
1920    expires_at_ms: Option<u64>,
1921) -> Result<(), String> {
1922    if let Some(expires_at_ms) = expires_at_ms {
1923        state
1924            .client()
1925            .register_token_until(token, scope, max_connections, expires_at_ms);
1926    } else {
1927        state
1928            .client()
1929            .register_session_token(token, scope, max_connections);
1930    }
1931    Ok(())
1932}
1933
1934#[tauri::command]
1935async fn get_endpoint_ticket_with_token(
1936    state: tauri::State<'_, OpenRtcTauriState>,
1937    scope: String,
1938    max_connections: u32,
1939) -> Result<String, String> {
1940    ensure_iroh_node(state.client()).await?;
1941    state
1942        .client()
1943        .endpoint_ticket_with_token(&scope, max_connections)
1944        .await
1945        .map_err(|error| error.to_string())
1946}
1947
1948#[tauri::command]
1949async fn validate_session_token(
1950    state: tauri::State<'_, OpenRtcTauriState>,
1951    token: String,
1952    connection_id: Option<String>,
1953) -> Result<String, String> {
1954    if let Some(connection_id) = connection_id
1955        .as_deref()
1956        .map(str::trim)
1957        .filter(|value| !value.is_empty())
1958    {
1959        state
1960            .client()
1961            .validate_connection_token(&token, connection_id, None)
1962            .await
1963            .map_err(|error| error.to_string())
1964    } else {
1965        state.client().validate_session_token(&token)
1966    }
1967}
1968
1969#[tauri::command]
1970async fn revoke_session_tokens_by_scope(
1971    state: tauri::State<'_, OpenRtcTauriState>,
1972    scope: String,
1973) -> Result<Vec<String>, String> {
1974    let affected = state.client().revoke_tokens_by_scope(&scope).await;
1975    if revokes_managed_session(&scope) {
1976        // The cached managed-session result contains the compound ticket whose
1977        // token was just revoked. Keeping that record would make the next
1978        // same-user/device start return a credential the host can no longer
1979        // validate. Invalidate only this lifecycle cache; the durable native
1980        // device identity remains intact so explicit reauthentication can
1981        // re-enroll the same installation with a freshly minted ticket.
1982        if let Some(previous) = state.managed_session.lock().await.take() {
1983            previous.owner_client.stop_presence_loop();
1984            previous.owner_client.stop_auto_connect();
1985            previous.owner_client.stop_external_auto_connect().await;
1986        }
1987        *state.presence_loop_active.lock().await = false;
1988    }
1989    Ok(affected)
1990}
1991
1992#[tauri::command]
1993async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1994    state.client().stop_presence_loop();
1995    *state.presence_loop_active.lock().await = false;
1996    if let Some(active) = state.managed_session.lock().await.as_mut() {
1997        active.result.presence_started = false;
1998    }
1999    stop_subscription_by_id(&state, "internal-presence-loop").await;
2000    Ok(())
2001}
2002
2003/// Register one public capability with the Rust-owned native root.
2004///
2005/// Owner: openrtc-tauri-plugin native root registry.
2006/// Consumers: OpenRTC native capability adapters.
2007/// Introduced: 2026-08-27.
2008#[tauri::command]
2009async fn register_rtc_capability(
2010    state: tauri::State<'_, OpenRtcTauriState>,
2011    capability_key: String,
2012    avenue_kind: String,
2013    avenue_id: String,
2014    sparse_fanout: Option<bool>,
2015    architecture: Option<openrtc::native::RoomArchitectureMode>,
2016) -> Result<(), String> {
2017    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2018    let avenue_kind = required_native_capability_part(avenue_kind, "avenue kind")?;
2019    let avenue_id = required_native_capability_part(avenue_id, "avenue id")?;
2020    if capability_key != format!("{avenue_kind}:{avenue_id}") {
2021        return Err("native capability key does not match its avenue".to_string());
2022    }
2023    state
2024        .capability_registry
2025        .lock()
2026        .await
2027        .register_with_architecture(
2028            capability_key,
2029            avenue_kind,
2030            avenue_id,
2031            sparse_fanout.unwrap_or(false),
2032            architecture,
2033        )
2034}
2035
2036#[cfg(feature = "managed-group-encryption")]
2037fn managed_group_state_path<R: Runtime>(
2038    app: &tauri::AppHandle<R>,
2039    state: &OpenRtcTauriState,
2040    capability_key: &str,
2041    device_id: &str,
2042) -> Result<PathBuf, String> {
2043    let mut digest = Sha256::new();
2044    digest.update(b"openrtc-tauri-managed-group-state-v1:");
2045    digest.update(state.client().app_tag().as_bytes());
2046    digest.update(device_id.as_bytes());
2047    digest.update(capability_key.as_bytes());
2048    let name = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(digest.finalize());
2049    Ok(app_data_dir(app, state)?
2050        .join("managed-room-state")
2051        .join(format!("{name}.bin")))
2052}
2053
2054#[cfg(feature = "managed-group-encryption")]
2055fn write_private_managed_group_state(path: &std::path::Path, bytes: &[u8]) -> Result<(), String> {
2056    if bytes.is_empty() || bytes.len() > 16 * 1024 * 1024 {
2057        return Err("managed room state exceeds its protected bound".to_string());
2058    }
2059    let parent = path
2060        .parent()
2061        .ok_or_else(|| "managed room state path has no parent".to_string())?;
2062    std::fs::create_dir_all(parent)
2063        .map_err(|error| format!("create managed room state directory: {error}"))?;
2064    #[cfg(unix)]
2065    {
2066        use std::os::unix::fs::PermissionsExt;
2067        std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700))
2068            .map_err(|error| format!("protect managed room state directory: {error}"))?;
2069    }
2070    let temporary = parent.join(format!(
2071        ".{}.{}.tmp",
2072        path.file_name()
2073            .and_then(|name| name.to_str())
2074            .unwrap_or("state"),
2075        uuid::Uuid::new_v4().simple(),
2076    ));
2077    std::fs::write(&temporary, bytes)
2078        .map_err(|error| format!("write managed room state: {error}"))?;
2079    #[cfg(unix)]
2080    {
2081        use std::os::unix::fs::PermissionsExt;
2082        std::fs::set_permissions(&temporary, std::fs::Permissions::from_mode(0o600))
2083            .map_err(|error| format!("protect managed room state: {error}"))?;
2084    }
2085    std::fs::rename(&temporary, path)
2086        .map_err(|error| format!("replace managed room state: {error}"))?;
2087    Ok(())
2088}
2089
2090#[cfg(feature = "managed-group-encryption")]
2091fn persist_managed_group_actions<R: Runtime>(
2092    app: &tauri::AppHandle<R>,
2093    state: &OpenRtcTauriState,
2094    capability_key: &str,
2095    device_id: &str,
2096    actions: Vec<serde_json::Value>,
2097) -> Result<Vec<serde_json::Value>, String> {
2098    let mut outbound = Vec::with_capacity(actions.len());
2099    for action in actions {
2100        if action.get("type").and_then(serde_json::Value::as_str) == Some("persist-state") {
2101            let sealed = action
2102                .get("sealedState")
2103                .and_then(serde_json::Value::as_str)
2104                .ok_or_else(|| "managed room persisted action is invalid".to_string())?;
2105            let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
2106                .decode(sealed)
2107                .map_err(|_| "managed room persisted action encoding is invalid".to_string())?;
2108            write_private_managed_group_state(
2109                &managed_group_state_path(app, state, capability_key, device_id)?,
2110                &bytes,
2111            )?;
2112        } else {
2113            outbound.push(action);
2114        }
2115    }
2116    Ok(outbound)
2117}
2118
2119#[cfg(feature = "managed-group-encryption")]
2120fn derive_tauri_managed_group_key(
2121    state: &OpenRtcTauriState,
2122    capability_key: &str,
2123    device_id: &str,
2124) -> Result<[u8; 32], String> {
2125    let client = state.client();
2126    let app_tag = client.app_tag();
2127    let challenge =
2128        format!("openrtc:managed-group-wrapping-key:v1:{app_tag}:{device_id}:{capability_key}",);
2129    let signer = state.native_device_key_signer.as_ref().ok_or_else(|| {
2130        "managed rooms require the host's sign-only native device key".to_string()
2131    })?;
2132    let mut signature = signer.sign(app_tag, challenge.as_bytes())?;
2133    if signature.len() != 64 {
2134        return Err("native device signer returned an invalid managed-room signature".to_string());
2135    }
2136    let mut digest = Sha256::new();
2137    digest.update(b"openrtc:managed-group-key-derivation:v1:");
2138    digest.update(challenge.as_bytes());
2139    digest.update(&signature);
2140    signature.fill(0);
2141    Ok(digest.finalize().into())
2142}
2143
2144#[cfg(feature = "managed-group-encryption")]
2145#[tauri::command]
2146async fn openrtc_init_managed_room_group<R: Runtime>(
2147    app: tauri::AppHandle<R>,
2148    state: tauri::State<'_, OpenRtcTauriState>,
2149    capability_key: String,
2150    device_id: String,
2151) -> Result<serde_json::Value, String> {
2152    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2153    let device_id = required_native_capability_part(device_id, "device id")?;
2154    {
2155        let registry = state.capability_registry.lock().await;
2156        let registration = registry
2157            .registrations
2158            .get(&capability_key)
2159            .ok_or_else(|| "managed room capability is not registered".to_string())?;
2160        if registration.avenue_kind != "room" {
2161            return Err("managed fanout is available only for room capabilities".to_string());
2162        }
2163    }
2164    let path = managed_group_state_path(&app, &state, &capability_key, &device_id)?;
2165    let sealed = match std::fs::read(path) {
2166        Ok(bytes) if bytes.len() <= 16 * 1024 * 1024 => Some(bytes),
2167        Ok(_) => return Err("managed room persisted state exceeds its protected bound".to_string()),
2168        Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
2169        Err(error) => return Err(format!("read managed room persisted state: {error}")),
2170    };
2171    let controller = openrtc::native::NativeManagedGroupController::new(
2172        &device_id,
2173        derive_tauri_managed_group_key(&state, &capability_key, &device_id)?,
2174        sealed.as_deref(),
2175    )
2176    .map_err(|error| format!("initialize managed room group: {error:#}"))?;
2177    let action = controller
2178        .publish_key_package()
2179        .map_err(|error| format!("create managed room KeyPackage: {error:#}"))?;
2180    state
2181        .managed_groups
2182        .lock()
2183        .await
2184        .insert(capability_key, controller);
2185    Ok(action)
2186}
2187
2188#[cfg(feature = "managed-group-encryption")]
2189#[tauri::command]
2190async fn openrtc_handle_managed_room_prepare_page<R: Runtime>(
2191    app: tauri::AppHandle<R>,
2192    state: tauri::State<'_, OpenRtcTauriState>,
2193    capability_key: String,
2194    device_id: String,
2195    page: serde_json::Value,
2196) -> Result<Vec<serde_json::Value>, String> {
2197    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2198    let device_id = required_native_capability_part(device_id, "device id")?;
2199    let actions = state
2200        .managed_groups
2201        .lock()
2202        .await
2203        .get_mut(&capability_key)
2204        .ok_or_else(|| "managed room group is not initialized".to_string())?
2205        .handle_prepare_page(page)
2206        .map_err(|error| format!("handle managed room preparation: {error:#}"))?;
2207    persist_managed_group_actions(&app, &state, &capability_key, &device_id, actions)
2208}
2209
2210#[cfg(feature = "managed-group-encryption")]
2211#[tauri::command]
2212async fn openrtc_handle_managed_room_artifact_chunk<R: Runtime>(
2213    app: tauri::AppHandle<R>,
2214    state: tauri::State<'_, OpenRtcTauriState>,
2215    capability_key: String,
2216    device_id: String,
2217    chunk: serde_json::Value,
2218) -> Result<Vec<serde_json::Value>, String> {
2219    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2220    let device_id = required_native_capability_part(device_id, "device id")?;
2221    let actions = state
2222        .managed_groups
2223        .lock()
2224        .await
2225        .get_mut(&capability_key)
2226        .ok_or_else(|| "managed room group is not initialized".to_string())?
2227        .handle_artifact_chunk(chunk)
2228        .map_err(|error| format!("handle managed room artifact: {error:#}"))?;
2229    persist_managed_group_actions(&app, &state, &capability_key, &device_id, actions)
2230}
2231
2232#[cfg(feature = "managed-group-encryption")]
2233#[tauri::command]
2234#[allow(clippy::too_many_arguments)]
2235async fn openrtc_seal_managed_room_payload<R: Runtime>(
2236    app: tauri::AppHandle<R>,
2237    state: tauri::State<'_, OpenRtcTauriState>,
2238    capability_key: String,
2239    device_id: String,
2240    architecture_epoch: u64,
2241    encryption_epoch: u64,
2242    message_id: String,
2243    channel: String,
2244    priority: u8,
2245    zone_id: Option<String>,
2246    payload: Vec<u8>,
2247) -> Result<openrtc::native::NativeManagedProtectedPayload, String> {
2248    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2249    let device_id = required_native_capability_part(device_id, "device id")?;
2250    let protected = state
2251        .managed_groups
2252        .lock()
2253        .await
2254        .get_mut(&capability_key)
2255        .ok_or_else(|| "managed room group is not initialized".to_string())?
2256        .seal_payload(
2257            architecture_epoch,
2258            encryption_epoch,
2259            &message_id,
2260            &channel,
2261            priority,
2262            zone_id.as_deref(),
2263            &payload,
2264        )
2265        .map_err(|error| format!("seal managed room payload: {error:#}"))?;
2266    let sealed = base64::engine::general_purpose::URL_SAFE_NO_PAD
2267        .decode(&protected.sealed_state)
2268        .map_err(|_| "managed room state encoding is invalid".to_string())?;
2269    write_private_managed_group_state(
2270        &managed_group_state_path(&app, &state, &capability_key, &device_id)?,
2271        &sealed,
2272    )?;
2273    Ok(protected)
2274}
2275
2276#[cfg(feature = "managed-group-encryption")]
2277#[tauri::command]
2278#[allow(clippy::too_many_arguments)]
2279async fn openrtc_open_managed_room_payload<R: Runtime>(
2280    app: tauri::AppHandle<R>,
2281    state: tauri::State<'_, OpenRtcTauriState>,
2282    capability_key: String,
2283    device_id: String,
2284    architecture_epoch: u64,
2285    encryption_epoch: u64,
2286    message_id: String,
2287    channel: String,
2288    priority: u8,
2289    zone_id: Option<String>,
2290    ciphertext: Vec<u8>,
2291) -> Result<openrtc::native::NativeManagedProtectedPayload, String> {
2292    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2293    let device_id = required_native_capability_part(device_id, "device id")?;
2294    let protected = state
2295        .managed_groups
2296        .lock()
2297        .await
2298        .get_mut(&capability_key)
2299        .ok_or_else(|| "managed room group is not initialized".to_string())?
2300        .open_payload(
2301            architecture_epoch,
2302            encryption_epoch,
2303            &message_id,
2304            &channel,
2305            priority,
2306            zone_id.as_deref(),
2307            &ciphertext,
2308        )
2309        .map_err(|error| format!("open managed room payload: {error:#}"))?;
2310    let sealed = base64::engine::general_purpose::URL_SAFE_NO_PAD
2311        .decode(&protected.sealed_state)
2312        .map_err(|_| "managed room state encoding is invalid".to_string())?;
2313    write_private_managed_group_state(
2314        &managed_group_state_path(&app, &state, &capability_key, &device_id)?,
2315        &sealed,
2316    )?;
2317    Ok(protected)
2318}
2319
2320#[cfg(feature = "managed-group-encryption")]
2321#[tauri::command]
2322async fn openrtc_forget_managed_room_group(
2323    state: tauri::State<'_, OpenRtcTauriState>,
2324    capability_key: String,
2325) -> Result<(), String> {
2326    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2327    state.managed_groups.lock().await.remove(&capability_key);
2328    Ok(())
2329}
2330
2331#[cfg(not(feature = "managed-group-encryption"))]
2332macro_rules! unavailable_managed_room_command {
2333    ($name:ident ( $($arg:ident : $type:ty),* ) -> $return:ty) => {
2334        #[tauri::command]
2335        async fn $name($($arg: $type),*) -> Result<$return, String> {
2336            $(let _ = $arg;)*
2337            Err("native managed rooms were not compiled into this Tauri host".to_string())
2338        }
2339    };
2340}
2341
2342#[cfg(not(feature = "managed-group-encryption"))]
2343unavailable_managed_room_command!(openrtc_init_managed_room_group(
2344    capability_key: String, device_id: String
2345) -> serde_json::Value);
2346#[cfg(not(feature = "managed-group-encryption"))]
2347unavailable_managed_room_command!(openrtc_handle_managed_room_prepare_page(
2348    capability_key: String, device_id: String, page: serde_json::Value
2349) -> Vec<serde_json::Value>);
2350#[cfg(not(feature = "managed-group-encryption"))]
2351unavailable_managed_room_command!(openrtc_handle_managed_room_artifact_chunk(
2352    capability_key: String, device_id: String, chunk: serde_json::Value
2353) -> Vec<serde_json::Value>);
2354#[cfg(not(feature = "managed-group-encryption"))]
2355unavailable_managed_room_command!(openrtc_seal_managed_room_payload(
2356    capability_key: String, device_id: String, architecture_epoch: u64,
2357    encryption_epoch: u64, message_id: String, channel: String, priority: u8,
2358    zone_id: Option<String>, payload: Vec<u8>
2359) -> serde_json::Value);
2360#[cfg(not(feature = "managed-group-encryption"))]
2361unavailable_managed_room_command!(openrtc_open_managed_room_payload(
2362    capability_key: String, device_id: String, architecture_epoch: u64,
2363    encryption_epoch: u64, message_id: String, channel: String, priority: u8,
2364    zone_id: Option<String>, ciphertext: Vec<u8>
2365) -> serde_json::Value);
2366#[cfg(not(feature = "managed-group-encryption"))]
2367unavailable_managed_room_command!(openrtc_forget_managed_room_group(
2368    capability_key: String
2369) -> ());
2370
2371#[tauri::command]
2372async fn unregister_rtc_capability(
2373    state: tauri::State<'_, OpenRtcTauriState>,
2374    capability_key: String,
2375) -> Result<(), String> {
2376    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2377    let (empty, aggregate) = {
2378        let mut registry = state.capability_registry.lock().await;
2379        registry.registrations.remove(&capability_key);
2380        let empty = registry.registrations.is_empty();
2381        let aggregate = (!empty)
2382            .then(|| registry.aggregate_desired_peers())
2383            .transpose()?;
2384        (empty, aggregate)
2385    };
2386    if empty {
2387        state.client().stop_external_auto_connect().await;
2388    } else if let Some((revision, peers_json)) = aggregate {
2389        state
2390            .client()
2391            .submit_external_desired_peers(revision, &peers_json)
2392            .await
2393            .map_err(|error| error.to_string())?;
2394    }
2395    Ok(())
2396}
2397
2398#[tauri::command]
2399async fn start_rtc_managed_session<R: Runtime>(
2400    app: tauri::AppHandle<R>,
2401    state: tauri::State<'_, OpenRtcTauriState>,
2402    user_id: String,
2403    device_name: Option<String>,
2404    local_device_id: Option<String>,
2405    metadata: Option<String>,
2406    transports: Option<openrtc::client::TransportConfig>,
2407    auto_connect: Option<bool>,
2408    presence: Option<bool>,
2409) -> Result<StartSessionResult, String> {
2410    let command_started = std::time::Instant::now();
2411    let _start_guard = state.managed_session_start_guard.lock().await;
2412    let client = state.client();
2413    let app_tag = client.app_tag().to_string();
2414    eprintln!(
2415        "[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
2416        user_id,
2417        app_tag,
2418        device_name
2419            .as_deref()
2420            .map(str::trim)
2421            .map(|value| !value.is_empty())
2422            .unwrap_or(false),
2423        local_device_id
2424            .as_deref()
2425            .map(str::trim)
2426            .map(|value| !value.is_empty())
2427            .unwrap_or(false),
2428        metadata
2429            .as_deref()
2430            .map(str::trim)
2431            .map(|value| !value.is_empty())
2432            .unwrap_or(false),
2433        auto_connect.unwrap_or(false),
2434        presence.unwrap_or(false)
2435    );
2436    let data_dir = app_data_dir(&app, &state)?;
2437    let install_context = InstallContext {
2438        data_dir: data_dir.clone(),
2439    };
2440    ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
2441        .await?;
2442    if let Some(transport_config) = transports {
2443        client
2444            .update_transport_config(transport_config)
2445            .await
2446            .map_err(|error| error.to_string())?;
2447    }
2448
2449    let local_node_id = ensure_iroh_node(client.clone()).await?;
2450    let local_device = client
2451        .init_native_device_identity(data_dir, device_name.as_deref())
2452        .await
2453        .map_err(|error| error.to_string())?;
2454
2455    let effective_local_device_id =
2456        managed_session_device_id(local_device_id.as_deref(), &local_device);
2457    if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
2458        if requested != local_device.device_id {
2459            eprintln!(
2460                "[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
2461                requested, local_device.device_id
2462            );
2463        }
2464    }
2465
2466    let should_start_presence = presence.unwrap_or(false);
2467    let should_start_auto_connect = auto_connect.unwrap_or(false);
2468    if should_start_presence || should_start_auto_connect {
2469        return Err(
2470            "native Firebase coordination has been removed; publish gateway presence and submit external desired peers instead"
2471                .to_string(),
2472        );
2473    }
2474    // The native endpoint and durable device are root-owned. Capability
2475    // principals live in the registry above and must not replace this physical
2476    // lifecycle owner when devices, spaces, rooms, or tickets coexist.
2477    let managed_session_key = format!("{}:{}", app_tag, effective_local_device_id);
2478    let (disposition, reused) = {
2479        let active = state.managed_session.lock().await;
2480        let disposition = managed_session_disposition(
2481            active.as_ref().map(|record| record.key.as_str()),
2482            active
2483                .as_ref()
2484                .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
2485            active
2486                .as_ref()
2487                .is_some_and(|record| record.result.presence_started),
2488            active
2489                .as_ref()
2490                .is_some_and(|record| record.result.auto_connect_started),
2491            &managed_session_key,
2492            should_start_presence,
2493            should_start_auto_connect,
2494        );
2495        let reused = (disposition == SessionDisposition::Reuse).then(|| {
2496            active
2497                .as_ref()
2498                .expect("managed session reuse requires an active record")
2499                .result
2500                .clone()
2501        });
2502        (disposition, reused)
2503    };
2504    if let Some(reused) = reused {
2505        eprintln!(
2506            "[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
2507            managed_session_key,
2508            command_started.elapsed().as_millis()
2509        );
2510        return Ok(reused);
2511    }
2512
2513    // A different app/user/device tuple replaces the previous lifecycle owner.
2514    // Stop old background work before publishing the new authoritative lease.
2515    if disposition == SessionDisposition::Replace {
2516        let previous = state
2517            .managed_session
2518            .lock()
2519            .await
2520            .take()
2521            .expect("managed session replacement requires an active record");
2522        previous.owner_client.stop_presence_loop();
2523        previous.owner_client.stop_auto_connect();
2524        previous.owner_client.stop_external_auto_connect().await;
2525        *state.presence_loop_active.lock().await = false;
2526    }
2527    let ticket_scope = Some("user-device".to_string());
2528    let ticket = client
2529        .endpoint_ticket_with_token("user-device", 0)
2530        .await
2531        .map_err(|error| error.to_string())?;
2532
2533    let logical_local_device = local_device.clone();
2534
2535    let owner_epoch = match disposition {
2536        SessionDisposition::Refresh => {
2537            state
2538                .managed_session
2539                .lock()
2540                .await
2541                .as_ref()
2542                .expect("managed session refresh requires an active record")
2543                .owner_epoch
2544        }
2545        SessionDisposition::Start | SessionDisposition::Replace => {
2546            state.allocate_managed_session_owner_epoch()
2547        }
2548        SessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
2549    };
2550
2551    eprintln!(
2552        "[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={}",
2553        user_id,
2554        app_tag,
2555        local_node_id,
2556        logical_local_device.device_id,
2557        owner_epoch,
2558        should_start_presence,
2559        should_start_auto_connect,
2560        ticket_scope.as_deref().unwrap_or("unrestricted"),
2561        command_started.elapsed().as_millis()
2562    );
2563
2564    let result = StartSessionResult {
2565        local_node_id,
2566        ticket_scope,
2567        ticket: Some(ticket),
2568        presence_started: false,
2569        auto_connect_started: false,
2570        local_device: logical_local_device,
2571    };
2572    *state.managed_session.lock().await = Some(ManagedSessionRecord {
2573        key: managed_session_key,
2574        owner_client: client,
2575        owner_epoch,
2576        result: result.clone(),
2577    });
2578    Ok(result)
2579}
2580
2581#[tauri::command]
2582async fn start_rtc_external_auto_connect(
2583    state: tauri::State<'_, OpenRtcTauriState>,
2584    user_id: String,
2585    local_device_id: String,
2586    capability_key: Option<String>,
2587) -> Result<(), String> {
2588    if let Some(capability_key) = capability_key {
2589        let capability_key = required_native_capability_part(capability_key, "capability key")?;
2590        if !state
2591            .capability_registry
2592            .lock()
2593            .await
2594            .registrations
2595            .contains_key(&capability_key)
2596        {
2597            return Err(format!(
2598                "native capability {capability_key} is not registered"
2599            ));
2600        }
2601    }
2602    // The external desired-peer actor is one Rust root actor. Public avenue
2603    // principals are registration metadata, not lifecycle keys.
2604    let root_user_id = format!("{}:native-root", state.client().app_tag());
2605    let _requested_capability_principal = user_id;
2606    state
2607        .client()
2608        .start_external_auto_connect(root_user_id, local_device_id)
2609        .await
2610        .map_err(|error| error.to_string())
2611}
2612
2613#[tauri::command]
2614async fn submit_rtc_desired_peers(
2615    state: tauri::State<'_, OpenRtcTauriState>,
2616    revision: u64,
2617    peers_json: String,
2618    capability_key: Option<String>,
2619    local_device_id: Option<String>,
2620    app_tag: Option<String>,
2621) -> Result<bool, String> {
2622    let original_peers = serde_json::from_str::<Vec<serde_json::Value>>(&peers_json)
2623        .map_err(|error| format!("invalid capability desired-peer payload: {error}"))?;
2624    if original_peers.len() > 100 {
2625        return Err("capability desired-peer payload exceeds 100 peers".to_string());
2626    }
2627    let key = {
2628        let registry = state.capability_registry.lock().await;
2629        match capability_key {
2630            Some(key) => required_native_capability_part(key, "capability key")?,
2631            None if registry.registrations.len() == 1 => registry
2632                .registrations
2633                .keys()
2634                .next()
2635                .cloned()
2636                .expect("one registration has one key"),
2637            None => {
2638                return Err(
2639                    "capabilityKey is required when multiple native capabilities are active"
2640                        .to_string(),
2641                )
2642            }
2643        }
2644    };
2645    let sparse = state
2646        .capability_registry
2647        .lock()
2648        .await
2649        .registrations
2650        .get(&key)
2651        .ok_or_else(|| format!("native capability {key} is not registered"))?
2652        .uses_sparse_fanout();
2653    let peers = if sparse {
2654        let local_device_id = local_device_id
2655            .as_deref()
2656            .map(str::trim)
2657            .filter(|value| !value.is_empty())
2658            .ok_or_else(|| "sparse native capability requires localDeviceId".to_string())?;
2659        let app_tag = app_tag
2660            .as_deref()
2661            .map(str::trim)
2662            .filter(|value| !value.is_empty())
2663            .ok_or_else(|| "sparse native capability requires appTag".to_string())?;
2664        let local_public_jwk = state.device_public_key(app_tag)?;
2665        let local_device_key_x = local_public_jwk
2666            .get("x")
2667            .and_then(serde_json::Value::as_str)
2668            .ok_or_else(|| "native sparse fanout signer returned no public key".to_string())?;
2669        let projected = state
2670            .client()
2671            .configure_sparse_fanout(
2672                &key,
2673                revision,
2674                local_device_id,
2675                local_device_key_x,
2676                &peers_json,
2677                true,
2678            )
2679            .await
2680            .map_err(|error| error.to_string())?
2681            .0;
2682        serde_json::from_str(&projected)
2683            .map_err(|error| format!("invalid projected desired-peer payload: {error}"))?
2684    } else {
2685        original_peers
2686    };
2687    let (root_revision, aggregate) = {
2688        let mut registry = state.capability_registry.lock().await;
2689        let registration = registry
2690            .registrations
2691            .get_mut(&key)
2692            .ok_or_else(|| format!("native capability {key} is not registered"))?;
2693        if revision <= registration.desired_revision {
2694            return Ok(false);
2695        }
2696        registration.desired_revision = revision;
2697        registration.desired_peers = peers;
2698        registry.aggregate_desired_peers()?
2699    };
2700    let accepted = state
2701        .client()
2702        .submit_external_desired_peers(root_revision, &aggregate)
2703        .await
2704        .map_err(|error| error.to_string())?;
2705    if accepted {
2706        // A capability can be registered after its peer is already connected.
2707        // Replay the Rust-owned read model after the desired-peer mapping is
2708        // committed instead of asking the frontend to redial or poll.
2709        replay_current_connection_states(&state).await;
2710    }
2711    Ok(accepted)
2712}
2713
2714#[tauri::command]
2715async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
2716    state.client().stop_auto_connect();
2717    state.client().stop_external_auto_connect().await;
2718    if let Some(active) = state.managed_session.lock().await.as_mut() {
2719        active.result.auto_connect_started = false;
2720    }
2721    Ok(())
2722}
2723
2724#[tauri::command]
2725async fn connect_to_device(
2726    state: tauri::State<'_, OpenRtcTauriState>,
2727    device_id: Option<String>,
2728    endpoint_ticket: String,
2729    timeout_ms: Option<u64>,
2730) -> Result<openrtc::client::ManagedConnectResult, String> {
2731    let client = state.client();
2732    let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
2733    match timeout_ms {
2734        Some(timeout_ms) => {
2735            tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
2736                .await
2737                .map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
2738                .map_err(|error| error.to_string())
2739        }
2740        None => connect.await.map_err(|error| error.to_string()),
2741    }
2742}
2743
2744#[tauri::command]
2745async fn disconnect_device(
2746    state: tauri::State<'_, OpenRtcTauriState>,
2747    device_id: String,
2748    node_id_hint: Option<String>,
2749) -> Result<(), String> {
2750    state
2751        .client()
2752        .disconnect_device(&device_id, node_id_hint.as_deref())
2753        .await;
2754    Ok(())
2755}
2756
2757#[tauri::command]
2758async fn set_auto_connect_excluded(
2759    state: tauri::State<'_, OpenRtcTauriState>,
2760    device_id: String,
2761    excluded: bool,
2762) -> Result<(), String> {
2763    if excluded {
2764        state.client().exclude_peer_and_publish(&device_id).await;
2765    } else {
2766        state.client().unexclude_peer_and_publish(&device_id).await;
2767    }
2768    Ok(())
2769}
2770
2771#[tauri::command]
2772async fn set_rtc_external_auto_connect_excluded(
2773    state: tauri::State<'_, OpenRtcTauriState>,
2774    device_id: String,
2775    excluded: bool,
2776) -> Result<(), String> {
2777    state
2778        .client()
2779        .set_auto_connect_excluded(&device_id, excluded);
2780    state.client().wake_native_external_auto_connect().await;
2781    Ok(())
2782}
2783
2784#[tauri::command]
2785async fn resolve_rtc_peer_connection_records(
2786    state: tauri::State<'_, OpenRtcTauriState>,
2787    id: String,
2788) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
2789    Ok(state.client().resolve_peer_connection_records(&id).await)
2790}
2791
2792#[tauri::command]
2793async fn resolve_rtc_peer_identity(
2794    state: tauri::State<'_, OpenRtcTauriState>,
2795    id: String,
2796) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
2797    Ok(state.client().peer_snapshot(&id).await)
2798}
2799
2800#[tauri::command]
2801async fn get_rtc_peer_session(
2802    state: tauri::State<'_, OpenRtcTauriState>,
2803    id: String,
2804) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2805    Ok(state.client().peer_session(&id).await)
2806}
2807
2808#[tauri::command]
2809async fn list_rtc_peer_sessions(
2810    state: tauri::State<'_, OpenRtcTauriState>,
2811) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
2812    Ok(state.client().peer_sessions().await)
2813}
2814
2815#[tauri::command]
2816async fn list_rtc_managed_connections(
2817    state: tauri::State<'_, OpenRtcTauriState>,
2818) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
2819    Ok(state.client().list_managed_connections().await)
2820}
2821
2822#[derive(Debug, Clone, Serialize)]
2823#[serde(rename_all = "camelCase")]
2824struct NotifyRtcNetworkChangeResult {
2825    retired_stale_connections: usize,
2826}
2827
2828#[tauri::command]
2829async fn notify_rtc_network_change(
2830    state: tauri::State<'_, OpenRtcTauriState>,
2831) -> Result<NotifyRtcNetworkChangeResult, String> {
2832    let retired_stale_connections = state
2833        .client()
2834        .notify_network_change()
2835        .await
2836        .map_err(|error| error.to_string())?;
2837    Ok(NotifyRtcNetworkChangeResult {
2838        retired_stale_connections,
2839    })
2840}
2841
2842#[tauri::command]
2843async fn set_rtc_transport_priority(
2844    state: tauri::State<'_, OpenRtcTauriState>,
2845    priority: Vec<openrtc::route_policy::KnownRoute>,
2846) -> Result<(), String> {
2847    let client = state.client();
2848    let mut transport_config = client.transport_config().await;
2849    transport_config.route_priority = priority;
2850    client
2851        .update_transport_config(transport_config)
2852        .await
2853        .map_err(|error| error.to_string())
2854}
2855
2856#[tauri::command]
2857async fn wait_for_rtc_settled_peer(
2858    state: tauri::State<'_, OpenRtcTauriState>,
2859    id: String,
2860    timeout_ms: Option<u64>,
2861) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2862    Ok(state.client().wait_for_peer(&id, timeout_ms).await)
2863}
2864
2865#[tauri::command]
2866async fn list_rtc_connection_states(
2867    state: tauri::State<'_, OpenRtcTauriState>,
2868    capability_key: String,
2869) -> Result<Vec<openrtc::client::StateSnapshot>, String> {
2870    let states = state.client().connection_states().await;
2871    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2872    let registry = state.capability_registry.lock().await;
2873    Ok(states
2874        .into_iter()
2875        .filter(|snapshot| {
2876            registry
2877                .capability_keys_for_state(snapshot)
2878                .iter()
2879                .any(|key| key == &capability_key)
2880        })
2881        .collect())
2882}
2883
2884#[tauri::command]
2885async fn get_rtc_connection_state(
2886    state: tauri::State<'_, OpenRtcTauriState>,
2887    connection_id: String,
2888    capability_key: String,
2889) -> Result<Option<openrtc::client::StateSnapshot>, String> {
2890    let snapshot = state.client().connection_state(&connection_id).await;
2891    let Some(snapshot) = snapshot else {
2892        return Ok(None);
2893    };
2894    let capability_key = required_native_capability_part(capability_key, "capability key")?;
2895    let included = state
2896        .capability_registry
2897        .lock()
2898        .await
2899        .capability_keys_for_state(&snapshot)
2900        .iter()
2901        .any(|key| key == &capability_key);
2902    Ok(included.then_some(snapshot))
2903}
2904
2905#[tauri::command]
2906async fn stop_rtc_subscription(
2907    state: tauri::State<'_, OpenRtcTauriState>,
2908    request_id: String,
2909) -> Result<(), String> {
2910    stop_subscription_by_id(&state, &request_id).await;
2911    Ok(())
2912}
2913
2914async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
2915    if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
2916        handle.abort();
2917    }
2918    let mut projected = state.native_projection_subscription.lock().await;
2919    if projected
2920        .as_ref()
2921        .is_some_and(|subscription| subscription.request_id == request_id)
2922    {
2923        *projected = None;
2924    }
2925}
2926
2927/// Register the one capability-keyed WebView projection Channel.
2928///
2929/// Rust forwards canonical connection state, decrypted peer data, and streams
2930/// that the embedding host explicitly hands to `handoff_projected_peer_bi_stream`.
2931/// The Channel is projection only; it never attaches another Rust receiver or
2932/// owns peer lifecycle.
2933#[tauri::command]
2934async fn start_native_projection(
2935    state: tauri::State<'_, OpenRtcTauriState>,
2936    request_id: Option<String>,
2937    channel: Channel<NativeProjectionEvent>,
2938) -> Result<String, String> {
2939    let request_id = request_id
2940        .map(|value| value.trim().to_string())
2941        .filter(|value| !value.is_empty())
2942        .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
2943    let mut subscription = state.native_projection_subscription.lock().await;
2944    install_native_projection(&mut subscription, request_id, channel)
2945}
2946
2947// The caller holds the existing projection-slot mutex across check and install.
2948fn install_native_projection(
2949    subscription: &mut Option<NativeProjectionSubscription>,
2950    request_id: String,
2951    channel: Channel<NativeProjectionEvent>,
2952) -> Result<String, String> {
2953    if let Some(current) = subscription.as_ref() {
2954        if current.request_id == request_id {
2955            // Same process owner reclaiming after Channel recreation (remount/HMR).
2956            *subscription = Some(NativeProjectionSubscription {
2957                request_id: request_id.clone(),
2958                channel,
2959            });
2960            return Ok(request_id);
2961        }
2962        // A different request is not proof of a stale owner. The current owner
2963        // must release its exact request before another Channel can claim it.
2964        return Err("OpenRTC native projection subscription already has a process owner".into());
2965    }
2966    *subscription = Some(NativeProjectionSubscription {
2967        request_id: request_id.clone(),
2968        channel,
2969    });
2970    Ok(request_id)
2971}
2972
2973#[tauri::command]
2974async fn get_projected_stream_diagnostics(
2975    state: tauri::State<'_, OpenRtcTauriState>,
2976) -> Result<ProjectedStreamDiagnostics, String> {
2977    Ok(ProjectedStreamDiagnostics {
2978        subscription_active: state.native_projection_subscription.lock().await.is_some(),
2979        offered: state.projected_stream_offered.load(Ordering::Relaxed),
2980        decoded: state.projected_stream_decoded.load(Ordering::Relaxed),
2981        projected: state.projected_stream_projected.load(Ordering::Relaxed),
2982        unhandled_non_channel: state
2983            .projected_stream_unhandled_non_channel
2984            .load(Ordering::Relaxed),
2985        unhandled_other_channel: state
2986            .projected_stream_unhandled_other_channel
2987            .load(Ordering::Relaxed),
2988        unhandled_no_subscription: state
2989            .projected_stream_unhandled_no_subscription
2990            .load(Ordering::Relaxed),
2991        failures: state.projected_stream_failures.load(Ordering::Relaxed),
2992        last_transport_stable_id: state
2993            .projected_stream_last_transport_stable_id
2994            .load(Ordering::Relaxed),
2995        validation_checks: state
2996            .projected_stream_validation_checks
2997            .load(Ordering::Relaxed),
2998        validation_rejections: state
2999            .projected_stream_validation_rejections
3000            .load(Ordering::Relaxed),
3001        last_validated_transport_stable_id: state
3002            .projected_stream_last_validated_transport_stable_id
3003            .load(Ordering::Relaxed),
3004        authorized: state.projected_stream_authorized.load(Ordering::Relaxed),
3005        unauthorized: state.projected_stream_unauthorized.load(Ordering::Relaxed),
3006        last_channel: state.projected_stream_last_channel.lock().await.clone(),
3007        last_protocol: state.projected_stream_last_protocol.lock().await.clone(),
3008        last_connection_id: state
3009            .projected_stream_last_connection_id
3010            .lock()
3011            .await
3012            .clone(),
3013        last_remote_node_id: state
3014            .projected_stream_last_remote_node_id
3015            .lock()
3016            .await
3017            .clone(),
3018    })
3019}
3020
3021#[tauri::command]
3022async fn is_current_transport_stable_id(
3023    state: tauri::State<'_, OpenRtcTauriState>,
3024    endpoint_id: String,
3025    transport_stable_id: u64,
3026) -> Result<bool, String> {
3027    state
3028        .projected_stream_validation_checks
3029        .fetch_add(1, Ordering::Relaxed);
3030    state
3031        .projected_stream_last_validated_transport_stable_id
3032        .store(transport_stable_id, Ordering::Relaxed);
3033    let is_current = state
3034        .client()
3035        .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
3036        .await
3037        .map_err(|error| error.to_string())?;
3038    if !is_current {
3039        state
3040            .projected_stream_validation_rejections
3041            .fetch_add(1, Ordering::Relaxed);
3042    }
3043    if native_stream_trace_enabled() {
3044        eprintln!(
3045            "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
3046            endpoint_id, transport_stable_id, is_current
3047        );
3048    }
3049    Ok(is_current)
3050}
3051
3052#[tauri::command]
3053async fn send_peer_message(
3054    state: tauri::State<'_, OpenRtcTauriState>,
3055    request: Request<'_>,
3056) -> Result<(), String> {
3057    let id = required_request_header(&request, PEER_ID_HEADER)?;
3058    let data = request_body_bytes(&request)?;
3059    state
3060        .client()
3061        .send_peer(&id, &data)
3062        .await
3063        .map_err(|error| error.to_string())
3064}
3065
3066#[tauri::command]
3067async fn encode_sparse_fanout_message(
3068    state: tauri::State<'_, OpenRtcTauriState>,
3069    request: Request<'_>,
3070) -> Result<Response, String> {
3071    let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
3072    let app_tag = required_request_header(&request, FANOUT_APP_TAG_HEADER)?;
3073    let payload = request_body_bytes(&request)?;
3074    let signing_request = state
3075        .client()
3076        .prepare_sparse_fanout_message(&capability, &payload)
3077        .await
3078        .map_err(|error| error.to_string())?;
3079    let signature = state.sign_device_message(&app_tag, &signing_request.signing_input)?;
3080    state
3081        .client()
3082        .finalize_sparse_fanout_message(&signing_request.request_id, &signature)
3083        .await
3084        .map(Response::new)
3085        .map_err(|error| error.to_string())
3086}
3087
3088#[tauri::command]
3089async fn accept_sparse_fanout_message(
3090    state: tauri::State<'_, OpenRtcTauriState>,
3091    request: Request<'_>,
3092) -> Result<serde_json::Value, String> {
3093    let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
3094    let source_peer = required_request_header(&request, FANOUT_SOURCE_PEER_HEADER)?;
3095    let encoded = request_body_bytes(&request)?;
3096    serde_json::to_value(
3097        state
3098            .client()
3099            .accept_sparse_fanout_message(&capability, &source_peer, &encoded)
3100            .await
3101            .map_err(|error| error.to_string())?,
3102    )
3103    .map_err(|error| error.to_string())
3104}
3105
3106#[tauri::command]
3107async fn sparse_fanout_diagnostics(
3108    state: tauri::State<'_, OpenRtcTauriState>,
3109    capability_key: String,
3110) -> Result<openrtc::sparse_fanout::SparseFanoutDiagnostics, String> {
3111    Ok(state
3112        .client()
3113        .sparse_fanout_diagnostics(&required_native_capability_part(
3114            capability_key,
3115            "capability key",
3116        )?)
3117        .await)
3118}
3119
3120#[derive(Debug, Serialize)]
3121#[serde(rename_all = "camelCase")]
3122struct NativeBroadcastPublisherDescriptor {
3123    handle: String,
3124    verification_key: Vec<u8>,
3125}
3126
3127#[derive(Debug, Serialize)]
3128#[serde(rename_all = "camelCase")]
3129struct NativeBroadcastSessionDescriptor {
3130    handle_id: String,
3131    id: String,
3132    role: openrtc::broadcast::BroadcastRole,
3133    state: openrtc::broadcast::BroadcastState,
3134    budget: openrtc::broadcast::BroadcastBudget,
3135}
3136
3137/// Creates a Rust-owned, sign-only publisher key. Only an opaque handle and
3138/// public verification key cross IPC; private bytes remain in this process.
3139#[tauri::command]
3140async fn prepare_openrtc_broadcast_publisher(
3141    state: tauri::State<'_, OpenRtcTauriState>,
3142) -> Result<NativeBroadcastPublisherDescriptor, String> {
3143    let signer = openrtc::broadcast::BroadcastPublisherSigner::generate()
3144        .map_err(|error| error.to_string())?;
3145    let verification_key = signer.verifying_key().to_vec();
3146    let handle = uuid::Uuid::new_v4().simple().to_string();
3147    state
3148        .broadcast_signers
3149        .lock()
3150        .await
3151        .insert(handle.clone(), signer);
3152    Ok(NativeBroadcastPublisherDescriptor {
3153        handle,
3154        verification_key,
3155    })
3156}
3157
3158#[tauri::command]
3159async fn release_openrtc_broadcast_publisher(
3160    state: tauri::State<'_, OpenRtcTauriState>,
3161    handle: String,
3162) -> Result<bool, String> {
3163    Ok(state
3164        .broadcast_signers
3165        .lock()
3166        .await
3167        .remove(&handle)
3168        .is_some())
3169}
3170
3171async fn drive_native_broadcast_actions(
3172    session: &openrtc::broadcast::BroadcastSession,
3173    adapter: &dyn NativeBroadcastAdapter,
3174) -> Result<(), String> {
3175    loop {
3176        let actions = session.take_actions(openrtc::broadcast::MAX_BROADCAST_ACTIONS);
3177        if actions.is_empty() {
3178            return Ok(());
3179        }
3180        for action in actions {
3181            if let Some(observation) = adapter.apply(session.clone(), action).await? {
3182                session
3183                    .observe(observation)
3184                    .map_err(|error| error.to_string())?;
3185            }
3186        }
3187    }
3188}
3189
3190async fn native_broadcast_driver(
3191    state: &OpenRtcTauriState,
3192    handle_id: &str,
3193) -> Result<Arc<Mutex<()>>, String> {
3194    state
3195        .broadcast_drivers
3196        .lock()
3197        .await
3198        .get(handle_id)
3199        .cloned()
3200        .ok_or_else(|| "broadcast session is unavailable".to_string())
3201}
3202
3203#[tauri::command]
3204async fn open_openrtc_broadcast(
3205    state: tauri::State<'_, OpenRtcTauriState>,
3206    grant_token: String,
3207    issuer_public_key: Vec<u8>,
3208    publisher_signer_handle: Option<String>,
3209    now_ms: u64,
3210) -> Result<NativeBroadcastSessionDescriptor, String> {
3211    let issuer_bytes: [u8; 32] = issuer_public_key
3212        .try_into()
3213        .map_err(|_| "broadcast issuer key must be 32 bytes".to_string())?;
3214    let issuer = VerifyingKey::from_bytes(&issuer_bytes)
3215        .map_err(|_| "broadcast issuer key is invalid".to_string())?;
3216    #[cfg(feature = "native-broadcast-moq")]
3217    let publication_verifying_key = if let Some(handle) = publisher_signer_handle.as_deref() {
3218        Some(
3219            state
3220                .broadcast_signers
3221                .lock()
3222                .await
3223                .get(handle)
3224                .ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?
3225                .verifying_key(),
3226        )
3227    } else {
3228        None
3229    };
3230    let broadcasts = state.client().broadcasts();
3231    let challenge = broadcasts
3232        .prepare_grant_verification(&grant_token, &issuer, now_ms)
3233        .map_err(|error| error.to_string())?;
3234    let device_signer = state.native_device_key_signer.as_deref().ok_or_else(|| {
3235        "native managed broadcast requires the host device-key signer".to_string()
3236    })?;
3237    let binding_signature = device_signer
3238        .sign(state.client().app_tag(), &challenge.signing_bytes)
3239        .map_err(|error| format!("sign broadcast installation challenge: {error}"))?;
3240    let grant = broadcasts
3241        .complete_grant_verification(
3242            &grant_token,
3243            &issuer,
3244            &challenge.handle,
3245            &binding_signature,
3246            now_ms,
3247        )
3248        .map_err(|error| error.to_string())?;
3249    let grant_generation = grant.grant_generation();
3250    #[cfg(feature = "native-broadcast-moq")]
3251    let native_moq_access = if state.native_broadcast_moq_adapter.is_some() {
3252        Some(
3253            state
3254                .resolve_native_broadcast_access(
3255                    &grant_token,
3256                    publication_verifying_key.as_ref(),
3257                    &issuer_bytes,
3258                )
3259                .await?,
3260        )
3261    } else {
3262        None
3263    };
3264    let session = if let Some(handle) = publisher_signer_handle.as_deref() {
3265        let signers = state.broadcast_signers.lock().await;
3266        let signer = signers
3267            .get(handle)
3268            .ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?;
3269        state
3270            .client()
3271            .broadcasts()
3272            .open_publisher(grant, signer, now_ms)
3273    } else {
3274        state.client().broadcasts().open(grant, now_ms)
3275    }
3276    .map_err(|error| error.to_string())?;
3277    let handle_id = format!("{}:{grant_generation}", session.id());
3278    #[cfg(feature = "native-broadcast-moq")]
3279    if let (Some(adapter), Some(access)) = (
3280        state.native_broadcast_moq_adapter.as_ref(),
3281        native_moq_access,
3282    ) {
3283        adapter
3284            .authorize(session.id(), grant_generation, access)
3285            .await;
3286    }
3287    let adapter = state.native_broadcast_adapter.as_deref().ok_or_else(|| {
3288        session.close();
3289        "native managed broadcast adapter is unavailable".to_string()
3290    })?;
3291    if let Err(error) = drive_native_broadcast_actions(&session, adapter).await {
3292        session.close();
3293        #[cfg(feature = "native-broadcast-moq")]
3294        if let Some(adapter) = state.native_broadcast_moq_adapter.as_ref() {
3295            adapter
3296                .forget_authorization(&session.id(), grant_generation)
3297                .await;
3298        }
3299        return Err(error);
3300    }
3301    let descriptor = NativeBroadcastSessionDescriptor {
3302        handle_id: handle_id.clone(),
3303        id: session.id(),
3304        role: session.role(),
3305        state: session.state(),
3306        budget: session.budget(),
3307    };
3308    state
3309        .broadcast_sessions
3310        .lock()
3311        .await
3312        .insert(handle_id.clone(), session);
3313    state
3314        .broadcast_drivers
3315        .lock()
3316        .await
3317        .insert(handle_id, Arc::new(Mutex::new(())));
3318    Ok(descriptor)
3319}
3320
3321#[tauri::command]
3322#[allow(clippy::too_many_arguments)]
3323async fn begin_openrtc_broadcast_publication(
3324    state: tauri::State<'_, OpenRtcTauriState>,
3325    handle_id: String,
3326    source_slot: String,
3327    publication_id: Option<String>,
3328    kind: String,
3329    codec: String,
3330    clock_rate: u32,
3331    coded_width: Option<u32>,
3332    coded_height: Option<u32>,
3333    channels: Option<u16>,
3334) -> Result<PortableMediaPublicationStart, String> {
3335    let driver = native_broadcast_driver(&state, &handle_id).await?;
3336    let _driver = driver.lock().await;
3337    let publication_id = match publication_id.as_deref().map(str::trim) {
3338        Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
3339        _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
3340    };
3341    let kind: openrtc::media::MediaKind = kind
3342        .parse()
3343        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3344    if !matches!(
3345        kind,
3346        openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
3347    ) {
3348        return Err("broadcast publications must be audio or video".to_string());
3349    }
3350    let publication = openrtc::media::MediaPublicationConfig {
3351        publication_id,
3352        media_generation: 0,
3353        kind,
3354        codec: codec
3355            .parse()
3356            .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?,
3357        clock_rate,
3358        coded_width,
3359        coded_height,
3360        channels,
3361    };
3362    let session = state
3363        .broadcast_sessions
3364        .lock()
3365        .await
3366        .get(&handle_id)
3367        .cloned()
3368        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3369    let publication = session
3370        .begin_publication(&source_slot, publication)
3371        .map_err(|error| error.to_string())?;
3372    let adapter = state
3373        .native_broadcast_adapter
3374        .as_deref()
3375        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3376    drive_native_broadcast_actions(&session, adapter).await?;
3377    Ok(PortableMediaPublicationStart {
3378        publication_id: publication.publication_id.to_string(),
3379        media_generation: publication.media_generation,
3380        // Broadcast control is signed and goes only to the private adapter.
3381        control: Vec::new(),
3382    })
3383}
3384
3385#[tauri::command]
3386#[allow(clippy::too_many_arguments)]
3387async fn publish_openrtc_broadcast_sample(
3388    state: tauri::State<'_, OpenRtcTauriState>,
3389    handle_id: String,
3390    publication_id: String,
3391    timestamp_us: u64,
3392    duration_us: u32,
3393    keyframe: bool,
3394    discardable: bool,
3395    payload: Vec<u8>,
3396) -> Result<(), String> {
3397    let driver = native_broadcast_driver(&state, &handle_id).await?;
3398    let _driver = driver.lock().await;
3399    openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
3400        .map_err(|error| error.to_string())?;
3401    let session = state
3402        .broadcast_sessions
3403        .lock()
3404        .await
3405        .get(&handle_id)
3406        .cloned()
3407        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3408    session
3409        .publish(
3410            parse_media_publication_id(&publication_id)?,
3411            openrtc::media::EncodedMediaSample {
3412                timestamp_us,
3413                duration_us,
3414                keyframe,
3415                discardable,
3416                payload,
3417            },
3418        )
3419        .map_err(|error| error.to_string())?;
3420    let adapter = state
3421        .native_broadcast_adapter
3422        .as_deref()
3423        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3424    drive_native_broadcast_actions(&session, adapter).await
3425}
3426
3427#[derive(Debug, Serialize)]
3428#[serde(rename_all = "camelCase")]
3429struct NativeBroadcastMediaSample {
3430    publication_id: String,
3431    media_generation: u32,
3432    sequence: u64,
3433    timestamp_us: u64,
3434    duration_us: u32,
3435    kind: &'static str,
3436    codec: &'static str,
3437    keyframe: bool,
3438    discardable: bool,
3439    payload: Vec<u8>,
3440}
3441
3442struct NativeBroadcastReceivedObject {
3443    source_slot: Option<String>,
3444    object: Vec<u8>,
3445}
3446
3447/// Retrieves only media objects already accepted by the private native
3448/// adapter and verified/decoded by the shared Rust broadcast owner. Provider
3449/// credentials, namespaces, and raw adapter actions never cross IPC.
3450#[tauri::command]
3451async fn receive_openrtc_broadcast_media(
3452    state: tauri::State<'_, OpenRtcTauriState>,
3453    handle_id: String,
3454    max: usize,
3455) -> Result<Vec<NativeBroadcastMediaSample>, String> {
3456    let driver = native_broadcast_driver(&state, &handle_id).await?;
3457    let _driver = driver.lock().await;
3458    let session = state
3459        .broadcast_sessions
3460        .lock()
3461        .await
3462        .get(&handle_id)
3463        .cloned()
3464        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3465    let adapter = state
3466        .native_broadcast_adapter
3467        .as_deref()
3468        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3469    drive_native_broadcast_actions(&session, adapter).await?;
3470    let objects: Vec<NativeBroadcastReceivedObject>;
3471    #[cfg(feature = "native-broadcast-moq")]
3472    {
3473        objects = if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
3474            native_moq
3475                .take_inbound_objects_with_source(&session, max.min(64))
3476                .await?
3477                .into_iter()
3478                .map(|item| NativeBroadcastReceivedObject {
3479                    source_slot: Some(item.source_slot),
3480                    object: item.object,
3481                })
3482                .collect()
3483        } else {
3484            adapter
3485                .take_inbound_objects(session.clone(), max.min(64))
3486                .await?
3487                .into_iter()
3488                .map(|object| NativeBroadcastReceivedObject {
3489                    source_slot: None,
3490                    object,
3491                })
3492                .collect()
3493        };
3494    }
3495    #[cfg(not(feature = "native-broadcast-moq"))]
3496    {
3497        objects = adapter
3498            .take_inbound_objects(session.clone(), max.min(64))
3499            .await?
3500            .into_iter()
3501            .map(|object| NativeBroadcastReceivedObject {
3502                source_slot: None,
3503                object,
3504            })
3505            .collect();
3506    }
3507    let mut accepted = Vec::with_capacity(objects.len());
3508    let mut delivered_bytes = 0_u64;
3509    for received in objects {
3510        let object_len = received.object.len() as u64;
3511        let chunk = if let Some(source_slot) = received.source_slot.as_deref() {
3512            let Some(media) = session
3513                .accept_media_for_source(source_slot, &received.object)
3514                .map_err(|error| error.to_string())?
3515            else {
3516                continue;
3517            };
3518            delivered_bytes = delivered_bytes.saturating_add(object_len);
3519            match media.media {
3520                Some(media) => Some(
3521                    openrtc::media::EncodedMediaChunk::decode(&media)
3522                        .map_err(|error| error.to_string())?,
3523                ),
3524                None => None,
3525            }
3526        } else {
3527            let chunk = session
3528                .accept_object(&received.object)
3529                .map_err(|error| error.to_string())?;
3530            delivered_bytes = delivered_bytes.saturating_add(object_len);
3531            chunk
3532        };
3533        let Some(chunk) = chunk else { continue };
3534        openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
3535            .and_then(|_| {
3536                openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
3537            })
3538            .map_err(|error| error.to_string())?;
3539        accepted.push(NativeBroadcastMediaSample {
3540            publication_id: chunk.publication_id.to_string(),
3541            media_generation: chunk.media_generation,
3542            sequence: chunk.sequence,
3543            timestamp_us: chunk.timestamp_us,
3544            duration_us: chunk.duration_us,
3545            kind: media_kind_name(chunk.kind),
3546            codec: media_codec_name(chunk.codec),
3547            keyframe: chunk.keyframe,
3548            discardable: chunk.discardable,
3549            payload: chunk.payload,
3550        });
3551    }
3552    #[cfg(feature = "native-broadcast-moq")]
3553    if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
3554        let usage = native_moq
3555            .record_delivered_bytes(&session, delivered_bytes)
3556            .await;
3557        let settlement = drive_native_broadcast_actions(&session, adapter).await;
3558        usage?;
3559        settlement?;
3560    }
3561    Ok(accepted)
3562}
3563
3564#[tauri::command]
3565async fn revoke_openrtc_broadcast(
3566    state: tauri::State<'_, OpenRtcTauriState>,
3567    handle_id: String,
3568    grant_generation: u64,
3569) -> Result<(), String> {
3570    let driver = native_broadcast_driver(&state, &handle_id).await?;
3571    let _driver = driver.lock().await;
3572    let session = state
3573        .broadcast_sessions
3574        .lock()
3575        .await
3576        .get(&handle_id)
3577        .cloned()
3578        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3579    session
3580        .revoke(grant_generation)
3581        .map_err(|error| error.to_string())?;
3582    let adapter = state
3583        .native_broadcast_adapter
3584        .as_deref()
3585        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3586    drive_native_broadcast_actions(&session, adapter).await
3587}
3588
3589#[tauri::command]
3590async fn close_openrtc_broadcast(
3591    state: tauri::State<'_, OpenRtcTauriState>,
3592    handle_id: String,
3593) -> Result<(), String> {
3594    let driver = native_broadcast_driver(&state, &handle_id).await?;
3595    let _driver = driver.lock().await;
3596    let session = state
3597        .broadcast_sessions
3598        .lock()
3599        .await
3600        .remove(&handle_id)
3601        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3602    session.close();
3603    let adapter = state
3604        .native_broadcast_adapter
3605        .as_deref()
3606        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3607    let result = drive_native_broadcast_actions(&session, adapter).await;
3608    state.broadcast_drivers.lock().await.remove(&handle_id);
3609    result
3610}
3611
3612#[tauri::command]
3613async fn record_sparse_fanout_forward_queue_drop(
3614    state: tauri::State<'_, OpenRtcTauriState>,
3615    capability_key: String,
3616    count: u64,
3617) -> Result<(), String> {
3618    state
3619        .client()
3620        .record_sparse_fanout_forward_queue_drop(
3621            &required_native_capability_part(capability_key, "capability key")?,
3622            count,
3623        )
3624        .await
3625        .map_err(|error| error.to_string())
3626}
3627
3628#[tauri::command]
3629async fn is_peer_connected(
3630    state: tauri::State<'_, OpenRtcTauriState>,
3631    node_id: String,
3632) -> Result<bool, String> {
3633    state
3634        .client()
3635        .is_connected_str(&node_id)
3636        .await
3637        .map_err(|error| error.to_string())
3638}
3639
3640#[derive(Debug, Serialize)]
3641#[serde(rename_all = "camelCase")]
3642struct PortableMediaPublicationStart {
3643    publication_id: String,
3644    media_generation: u32,
3645    control: Vec<u8>,
3646}
3647
3648fn media_kind_name(kind: openrtc::media::MediaKind) -> &'static str {
3649    match kind {
3650        openrtc::media::MediaKind::Audio => "audio",
3651        openrtc::media::MediaKind::Video => "video",
3652        openrtc::media::MediaKind::Screen => "screen",
3653        openrtc::media::MediaKind::Data => "data",
3654    }
3655}
3656
3657fn media_codec_name(codec: openrtc::media::MediaCodec) -> &'static str {
3658    match codec {
3659        openrtc::media::MediaCodec::Opus => "opus",
3660        openrtc::media::MediaCodec::H264 => "h264",
3661        openrtc::media::MediaCodec::Vp8 => "vp8",
3662        openrtc::media::MediaCodec::Vp9 => "vp9",
3663        openrtc::media::MediaCodec::Av1 => "av1",
3664        openrtc::media::MediaCodec::Pcm => "pcm",
3665        openrtc::media::MediaCodec::Opaque => "opaque",
3666    }
3667}
3668
3669fn parse_media_publication_id(value: &str) -> Result<openrtc::media::PublicationId, String> {
3670    value
3671        .trim()
3672        .parse()
3673        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())
3674}
3675
3676#[tauri::command]
3677#[allow(clippy::too_many_arguments)]
3678async fn begin_openrtc_media_publication(
3679    state: tauri::State<'_, OpenRtcTauriState>,
3680    publication_id: Option<String>,
3681    kind: String,
3682    codec: String,
3683    clock_rate: u32,
3684    coded_width: Option<u32>,
3685    coded_height: Option<u32>,
3686    channels: Option<u16>,
3687) -> Result<PortableMediaPublicationStart, String> {
3688    let publication_id = match publication_id.as_deref().map(str::trim) {
3689        Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
3690        _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
3691    };
3692    let kind: openrtc::media::MediaKind = kind
3693        .parse()
3694        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3695    if !matches!(
3696        kind,
3697        openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
3698    ) {
3699        return Err("browser media publications must be audio or video".to_string());
3700    }
3701    let codec = codec
3702        .parse()
3703        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3704    let publication = openrtc::media::MediaPublicationConfig {
3705        publication_id,
3706        media_generation: 0,
3707        kind,
3708        codec,
3709        clock_rate,
3710        coded_width,
3711        coded_height,
3712        channels,
3713    };
3714    let (publication, control) = state
3715        .portable_media
3716        .lock()
3717        .await
3718        .begin_publication(publication)
3719        .map_err(|error| error.to_string())?;
3720    Ok(PortableMediaPublicationStart {
3721        publication_id: publication.publication_id.to_string(),
3722        media_generation: publication.media_generation,
3723        control,
3724    })
3725}
3726
3727#[tauri::command]
3728#[allow(clippy::too_many_arguments)]
3729async fn encode_openrtc_media_sample(
3730    state: tauri::State<'_, OpenRtcTauriState>,
3731    request: Request<'_>,
3732) -> Result<Response, String> {
3733    let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
3734    let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
3735    openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
3736        .map_err(|error| error.to_string())?;
3737    let duration_us = parsed_request_header(&request, MEDIA_DURATION_US_HEADER)?;
3738    let keyframe = parsed_request_header(&request, MEDIA_KEYFRAME_HEADER)?;
3739    let discardable = parsed_request_header(&request, MEDIA_DISCARDABLE_HEADER)?;
3740    let payload = request_body_bytes(&request)?;
3741    let encoded = state
3742        .portable_media
3743        .lock()
3744        .await
3745        .encode_sample(
3746            parse_media_publication_id(&publication_id)?,
3747            openrtc::media::EncodedMediaSample {
3748                timestamp_us,
3749                duration_us,
3750                keyframe,
3751                discardable,
3752                payload,
3753            },
3754        )
3755        .map_err(|error| error.to_string())?;
3756    Ok(Response::new(encoded))
3757}
3758
3759#[tauri::command]
3760async fn pause_openrtc_media_publication(
3761    state: tauri::State<'_, OpenRtcTauriState>,
3762    publication_id: String,
3763) -> Result<(), String> {
3764    state
3765        .portable_media
3766        .lock()
3767        .await
3768        .pause_publication(parse_media_publication_id(&publication_id)?)
3769        .map_err(|error| error.to_string())
3770}
3771
3772#[tauri::command]
3773async fn retire_openrtc_media_publication(
3774    state: tauri::State<'_, OpenRtcTauriState>,
3775    publication_id: String,
3776) -> Result<(), String> {
3777    state
3778        .portable_media
3779        .lock()
3780        .await
3781        .retire_publication(parse_media_publication_id(&publication_id)?);
3782    Ok(())
3783}
3784
3785#[tauri::command]
3786async fn retire_openrtc_media_receiver(
3787    state: tauri::State<'_, OpenRtcTauriState>,
3788    publication_id: String,
3789    media_generation: u32,
3790) -> Result<(), String> {
3791    state.portable_media.lock().await.retire_receiver(
3792        parse_media_publication_id(&publication_id)?,
3793        media_generation,
3794    );
3795    Ok(())
3796}
3797
3798#[tauri::command]
3799async fn decode_openrtc_media_chunk(
3800    state: tauri::State<'_, OpenRtcTauriState>,
3801    request: Request<'_>,
3802) -> Result<Response, String> {
3803    let encoded = request_body_bytes(&request)?;
3804    let chunk = state
3805        .portable_media
3806        .lock()
3807        .await
3808        .decode_chunk(&encoded)
3809        .map_err(|error| error.to_string())?;
3810    openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
3811        .and_then(|_| {
3812            openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
3813        })
3814        .map_err(|error| error.to_string())?;
3815    let metadata = serde_json::to_vec(&serde_json::json!({
3816        "publicationId": chunk.publication_id.to_string(),
3817        "mediaGeneration": chunk.media_generation,
3818        "sequence": chunk.sequence,
3819        "timestampUs": chunk.timestamp_us,
3820        "durationUs": chunk.duration_us,
3821        "kind": media_kind_name(chunk.kind),
3822        "codec": media_codec_name(chunk.codec),
3823        "keyframe": chunk.keyframe,
3824        "discardable": chunk.discardable,
3825    }))
3826    .map_err(|error| format!("encode media IPC metadata failed: {error}"))?;
3827    let metadata_len =
3828        u32::try_from(metadata.len()).map_err(|_| "media IPC metadata is too large".to_string())?;
3829    let mut response = Vec::with_capacity(4 + metadata.len() + chunk.payload.len());
3830    response.extend_from_slice(&metadata_len.to_be_bytes());
3831    response.extend_from_slice(&metadata);
3832    response.extend_from_slice(&chunk.payload);
3833    Ok(Response::new(response))
3834}
3835
3836#[tauri::command]
3837async fn decode_openrtc_media_control(
3838    state: tauri::State<'_, OpenRtcTauriState>,
3839    request: Request<'_>,
3840) -> Result<serde_json::Value, String> {
3841    let encoded = request_body_bytes(&request)?;
3842    let control = state
3843        .portable_media
3844        .lock()
3845        .await
3846        .decode_control(&encoded)
3847        .map_err(|error| error.to_string())?;
3848    Ok(match control {
3849        openrtc::media::MediaControlFrame::Publish {
3850            publication_id,
3851            media_generation,
3852            kind,
3853            codec,
3854            clock_rate,
3855            coded_width,
3856            coded_height,
3857            channels,
3858        } => serde_json::json!({
3859            "type": "publish",
3860            "publicationId": publication_id.to_string(),
3861            "mediaGeneration": media_generation,
3862            "kind": media_kind_name(kind),
3863            "codec": media_codec_name(codec),
3864            "clockRate": clock_rate,
3865            "codedWidth": coded_width,
3866            "codedHeight": coded_height,
3867            "channels": channels,
3868        }),
3869        openrtc::media::MediaControlFrame::SetEnabled {
3870            publication_id,
3871            media_generation,
3872            enabled,
3873        } => serde_json::json!({
3874            "type": "set-enabled",
3875            "publicationId": publication_id.to_string(),
3876            "mediaGeneration": media_generation,
3877            "enabled": enabled,
3878        }),
3879        openrtc::media::MediaControlFrame::RequestKeyframe {
3880            publication_id,
3881            media_generation,
3882        } => serde_json::json!({
3883            "type": "request-keyframe",
3884            "publicationId": publication_id.to_string(),
3885            "mediaGeneration": media_generation,
3886        }),
3887        openrtc::media::MediaControlFrame::Stop {
3888            publication_id,
3889            media_generation,
3890            reason,
3891        } => serde_json::json!({
3892            "type": "stop",
3893            "publicationId": publication_id.to_string(),
3894            "mediaGeneration": media_generation,
3895            "reason": reason,
3896        }),
3897    })
3898}
3899
3900#[tauri::command]
3901async fn open_peer_bi_stream<R: Runtime>(
3902    app: tauri::AppHandle<R>,
3903    state: tauri::State<'_, OpenRtcTauriState>,
3904    peer_id: String,
3905    timeout_ms: Option<u64>,
3906) -> Result<OpenBiResult, String> {
3907    open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
3908        client.open_peer_bi(&peer_id, timeout_ms).await
3909    })
3910    .await
3911}
3912
3913#[tauri::command]
3914async fn open_peer_bi_transport_only_stream<R: Runtime>(
3915    app: tauri::AppHandle<R>,
3916    state: tauri::State<'_, OpenRtcTauriState>,
3917    peer_id: String,
3918    timeout_ms: Option<u64>,
3919) -> Result<OpenBiResult, String> {
3920    open_peer_bi_with(
3921        app,
3922        state,
3923        peer_id,
3924        timeout_ms,
3925        |client, peer_id, timeout_ms| async move {
3926            client
3927                .open_peer_bi_transport_only(&peer_id, timeout_ms)
3928                .await
3929                .map(|(connection_id, remote_node_id, send, recv)| {
3930                    (
3931                        connection_id,
3932                        remote_node_id,
3933                        PeerSendStream::plain(send),
3934                        PeerRecvStream::plain(recv),
3935                    )
3936                })
3937        },
3938    )
3939    .await
3940}
3941
3942#[tauri::command]
3943async fn open_peer_uni_stream(
3944    state: tauri::State<'_, OpenRtcTauriState>,
3945    peer_id: String,
3946    timeout_ms: Option<u64>,
3947) -> Result<OpenUniResult, String> {
3948    let peer_id = peer_id.trim().to_string();
3949    if peer_id.is_empty() {
3950        return Err("peerId is required".to_string());
3951    }
3952    let (connection_id, remote_node_id, send) = state
3953        .client()
3954        .open_peer_uni(&peer_id, timeout_ms)
3955        .await
3956        .map_err(|error| format!("open peer uni stream failed: {error}"))?;
3957    let stream_id = uuid::Uuid::new_v4().to_string();
3958    state
3959        .peer_uni_streams
3960        .lock()
3961        .await
3962        .insert(stream_id.clone(), Arc::new(Mutex::new(Some(send))));
3963    Ok(OpenUniResult {
3964        stream_id,
3965        connection_id,
3966        remote_node_id,
3967    })
3968}
3969
3970#[tauri::command]
3971async fn write_peer_bi_stream(
3972    state: tauri::State<'_, OpenRtcTauriState>,
3973    request: Request<'_>,
3974) -> Result<(), String> {
3975    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
3976    let bytes = request_body_bytes(&request)?;
3977
3978    let send = {
3979        let streams = state.peer_bi_streams.lock().await;
3980        streams
3981            .get(&stream_id)
3982            .map(|handle| handle.send.clone())
3983            .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
3984    };
3985    let result = send
3986        .lock()
3987        .await
3988        .as_mut()
3989        .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
3990        .write_all(&bytes)
3991        .await
3992        .map_err(|error| format!("write peer bi stream failed: {error}"));
3993    result
3994}
3995
3996#[tauri::command]
3997async fn start_peer_bi_stream_read<R: Runtime>(
3998    _app: tauri::AppHandle<R>,
3999    state: tauri::State<'_, OpenRtcTauriState>,
4000    stream_id: String,
4001    channel: Channel<Response>,
4002) -> Result<(), String> {
4003    let stream_id = stream_id.trim().to_string();
4004    if stream_id.is_empty() {
4005        return Err("streamId is required".to_string());
4006    }
4007
4008    let mut streams = state.peer_bi_streams.lock().await;
4009    let handle = streams
4010        .get_mut(&stream_id)
4011        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
4012    if handle.read_task.is_some() {
4013        return Ok(());
4014    }
4015    let recv = handle
4016        .recv
4017        .take()
4018        .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
4019    handle.read_task = Some(spawn_peer_bi_stream_reader(channel, stream_id, recv));
4020    Ok(())
4021}
4022
4023/// Cancel only the receive half of a bidirectional stream.
4024///
4025/// Web Streams permits a consumer to cancel an unused readable while keeping
4026/// its writable alive. Media publications rely on that half-close behavior, so
4027/// the native IPC projection must not remove or finish the Rust send handle.
4028#[tauri::command]
4029async fn cancel_peer_bi_stream_read(
4030    state: tauri::State<'_, OpenRtcTauriState>,
4031    stream_id: String,
4032) -> Result<(), String> {
4033    let stream_id = stream_id.trim().to_string();
4034    if stream_id.is_empty() {
4035        return Err("streamId is required".to_string());
4036    }
4037
4038    let mut streams = state.peer_bi_streams.lock().await;
4039    let handle = streams
4040        .get_mut(&stream_id)
4041        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
4042    handle.recv.take();
4043    if let Some(read_task) = handle.read_task.take() {
4044        read_task.abort();
4045    }
4046    Ok(())
4047}
4048
4049#[tauri::command]
4050async fn close_peer_bi_stream(
4051    state: tauri::State<'_, OpenRtcTauriState>,
4052    stream_id: String,
4053) -> Result<(), String> {
4054    let stream_id = stream_id.trim().to_string();
4055    if stream_id.is_empty() {
4056        return Err("streamId is required".to_string());
4057    }
4058
4059    let handle = {
4060        let mut streams = state.peer_bi_streams.lock().await;
4061        streams.remove(&stream_id)
4062    };
4063    if let Some(handle) = handle {
4064        // Hold the writer mutex before finishing so close waits for any
4065        // in-flight chunk write instead of skipping the flush when Arc clones exist.
4066        let mut send = handle.send.lock().await;
4067        if let Some(send) = send.take() {
4068            let _ = send.finish();
4069        }
4070        drop(send);
4071        if let Some(read_task) = handle.read_task {
4072            read_task.abort();
4073        }
4074    }
4075    Ok(())
4076}
4077
4078#[tauri::command]
4079async fn write_peer_uni_stream(
4080    state: tauri::State<'_, OpenRtcTauriState>,
4081    request: Request<'_>,
4082) -> Result<(), String> {
4083    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
4084    let bytes = request_body_bytes(&request)?;
4085    let send = {
4086        let streams = state.peer_uni_streams.lock().await;
4087        streams
4088            .get(&stream_id)
4089            .cloned()
4090            .ok_or_else(|| format!("peer uni stream not found: {stream_id}"))?
4091    };
4092    let result = send
4093        .lock()
4094        .await
4095        .as_mut()
4096        .ok_or_else(|| format!("peer uni stream already closed: {stream_id}"))?
4097        .write_all(&bytes)
4098        .await
4099        .map_err(|error| format!("write peer uni stream failed: {error}"));
4100    result
4101}
4102
4103#[tauri::command]
4104async fn close_peer_uni_stream(
4105    state: tauri::State<'_, OpenRtcTauriState>,
4106    stream_id: String,
4107) -> Result<(), String> {
4108    let stream_id = stream_id.trim().to_string();
4109    if stream_id.is_empty() {
4110        return Err("streamId is required".to_string());
4111    }
4112    let send = state.peer_uni_streams.lock().await.remove(&stream_id);
4113    if let Some(send) = send {
4114        if let Some(send) = send.lock().await.take() {
4115            let _ = send.finish();
4116        }
4117    }
4118    Ok(())
4119}
4120
4121#[cfg(test)]
4122mod tests {
4123    use super::*;
4124    use std::collections::HashMap as StdHashMap;
4125    use std::sync::atomic::{AtomicUsize, Ordering};
4126    use std::sync::Mutex as StdMutex;
4127    use tokio::sync::oneshot;
4128
4129    #[test]
4130    fn native_projection_rejects_competing_owner_without_replacing_channel() {
4131        let mut slot = None;
4132        let first = Channel::new(|_| Ok(()));
4133        let first_id = first.id();
4134        install_native_projection(&mut slot, "first".into(), first).unwrap();
4135        let error = install_native_projection(
4136            &mut slot, "second".into(), Channel::new(|_| Ok(())),
4137        ).unwrap_err();
4138        assert!(error.contains("already has a process owner"));
4139        assert_eq!(slot.as_ref().unwrap().request_id, "first");
4140        assert_eq!(slot.as_ref().unwrap().channel.id(), first_id);
4141
4142        let refreshed = Channel::new(|_| Ok(()));
4143        let refreshed_id = refreshed.id();
4144        install_native_projection(&mut slot, "first".into(), refreshed).unwrap();
4145        assert_eq!(slot.as_ref().unwrap().channel.id(), refreshed_id);
4146        slot = None;
4147        install_native_projection(&mut slot, "second".into(), Channel::new(|_| Ok(()))).unwrap();
4148        assert_eq!(slot.as_ref().unwrap().request_id, "second");
4149        assert!(install_native_projection(
4150            &mut slot, "first".into(), Channel::new(|_| Ok(())),
4151        ).is_err());
4152        assert_eq!(slot.as_ref().unwrap().request_id, "second");
4153    }
4154
4155    struct FakeTransportInstaller {
4156        installs: Arc<AtomicUsize>,
4157    }
4158
4159    #[derive(Default)]
4160    struct FakeDeviceKeySigner {
4161        records: StdMutex<StdHashMap<String, String>>,
4162    }
4163
4164    impl DeviceKeySigner for FakeDeviceKeySigner {
4165        fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
4166            Ok(serde_json::json!({
4167                "kty": "OKP",
4168                "crv": "Ed25519",
4169                "x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
4170            }))
4171        }
4172
4173        fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
4174            assert_eq!(challenge, b"openrtc:v2:test");
4175            Ok(vec![7; 64])
4176        }
4177
4178        fn delete(&self, _app_tag: &str) -> Result<(), String> {
4179            Ok(())
4180        }
4181
4182        fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
4183            Ok(self
4184                .records
4185                .lock()
4186                .unwrap()
4187                .get(&format!("{app_tag}:{key}"))
4188                .cloned())
4189        }
4190
4191        fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
4192            self.records
4193                .lock()
4194                .unwrap()
4195                .insert(format!("{app_tag}:{key}"), value.to_string());
4196            Ok(())
4197        }
4198
4199        fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
4200            self.records
4201                .lock()
4202                .unwrap()
4203                .remove(&format!("{app_tag}:{key}"));
4204            Ok(())
4205        }
4206    }
4207
4208    impl TransportInstaller for FakeTransportInstaller {
4209        fn id(&self) -> &'static str {
4210            "ble"
4211        }
4212
4213        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
4214            config.ble.as_ref().is_some_and(|ble| ble.enabled)
4215        }
4216
4217        fn install(
4218            &self,
4219            _client: Arc<openrtc::client::Client>,
4220            _config: openrtc::client::TransportConfig,
4221            _context: InstallContext,
4222        ) -> InstallFuture {
4223            self.installs.fetch_add(1, Ordering::SeqCst);
4224            Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
4225        }
4226    }
4227
4228    #[derive(Debug)]
4229    struct UnavailableTransportInstaller;
4230
4231    impl TransportInstaller for UnavailableTransportInstaller {
4232        fn id(&self) -> &'static str {
4233            "ble"
4234        }
4235
4236        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
4237            config.ble.as_ref().is_some_and(|ble| ble.enabled)
4238        }
4239
4240        fn install(
4241            &self,
4242            _client: Arc<openrtc::client::Client>,
4243            _config: openrtc::client::TransportConfig,
4244            _context: InstallContext,
4245        ) -> InstallFuture {
4246            Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
4247        }
4248    }
4249
4250    fn requested_ble_config() -> openrtc::client::TransportConfig {
4251        openrtc::client::TransportConfig {
4252            ble: Some(openrtc::client::BleConfig {
4253                enabled: true,
4254                ..openrtc::client::BleConfig::default()
4255            }),
4256            ..openrtc::client::TransportConfig::default()
4257        }
4258    }
4259
4260    fn test_config() -> OpenRtcTauriConfig {
4261        OpenRtcTauriConfig {
4262            api_key: format!("pk_test_{}", "a".repeat(40)),
4263            ..OpenRtcTauriConfig::default()
4264        }
4265    }
4266
4267    #[test]
4268    fn native_room_architecture_requires_fixed_intent_until_rust_owns_the_gateway_handle() {
4269        let mut registry = NativeCapabilityRegistry::default();
4270        registry
4271            .register_with_architecture(
4272                "room:adaptive".to_string(),
4273                "room".to_string(),
4274                "adaptive".to_string(),
4275                true,
4276                Some(openrtc::native::RoomArchitectureMode::Auto),
4277            )
4278            .expect("room registration");
4279        assert!(
4280            !registry.registrations["room:adaptive"].uses_sparse_fanout(),
4281            "caller-supplied auto state cannot promote sparse fanout",
4282        );
4283        registry
4284            .register_with_architecture(
4285                "room:fixed-sparse".to_string(),
4286                "room".to_string(),
4287                "fixed-sparse".to_string(),
4288                false,
4289                Some(openrtc::native::RoomArchitectureMode::Sparse),
4290            )
4291            .expect("fixed sparse room registration");
4292        assert!(registry.registrations["room:fixed-sparse"].uses_sparse_fanout());
4293    }
4294
4295    #[test]
4296    fn native_room_architecture_rejects_non_room_and_respects_fixed_mesh_intent() {
4297        let mut registry = NativeCapabilityRegistry::default();
4298        assert!(registry
4299            .register_with_architecture(
4300                "space:not-room".to_string(),
4301                "space".to_string(),
4302                "not-room".to_string(),
4303                false,
4304                Some(openrtc::native::RoomArchitectureMode::Auto),
4305            )
4306            .is_err());
4307        registry
4308            .register_with_architecture(
4309                "room:fixed".to_string(),
4310                "room".to_string(),
4311                "fixed".to_string(),
4312                false,
4313                Some(openrtc::native::RoomArchitectureMode::Mesh),
4314            )
4315            .expect("fixed room registration");
4316        assert!(!registry.registrations["room:fixed"].uses_sparse_fanout());
4317    }
4318
4319    #[test]
4320    fn native_capability_registry_isolates_duplicate_peer_channels() {
4321        let mut registry = NativeCapabilityRegistry::default();
4322        registry
4323            .register(
4324                "devices:user-1".to_string(),
4325                "devices".to_string(),
4326                "user-1".to_string(),
4327                false,
4328            )
4329            .expect("devices registration");
4330        registry
4331            .register(
4332                "space:room-1".to_string(),
4333                "space".to_string(),
4334                "room-1".to_string(),
4335                false,
4336            )
4337            .expect("space registration");
4338        registry
4339            .registrations
4340            .get_mut("devices:user-1")
4341            .unwrap()
4342            .desired_peers = vec![serde_json::json!({
4343            "deviceId": "shared-peer",
4344            "nodeId": "device-route"
4345        })];
4346        registry
4347            .registrations
4348            .get_mut("space:room-1")
4349            .unwrap()
4350            .desired_peers = vec![serde_json::json!({
4351            "deviceId": "shared-peer",
4352            "nodeId": "space-route"
4353        })];
4354
4355        assert_eq!(
4356            registry.capability_keys_for_identity(
4357                Some("shared-peer"),
4358                Some("shared-peer"),
4359                None,
4360                None,
4361            ),
4362            vec!["devices:user-1".to_string(), "space:room-1".to_string()]
4363        );
4364
4365        let explicit_channel = openrtc::stream_metadata::ChannelMetadata {
4366            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4367            metadata: Some(serde_json::Map::from_iter([(
4368                "openrtcCapability".to_string(),
4369                serde_json::json!("space:room-1"),
4370            )])),
4371        };
4372        assert_eq!(
4373            registry.projected_stream_capability(&explicit_channel),
4374            Some("space:room-1".to_string())
4375        );
4376
4377        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
4378            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4379            metadata: None,
4380        };
4381        assert_eq!(
4382            registry.projected_stream_capability(&unscoped_channel),
4383            None
4384        );
4385    }
4386
4387    #[test]
4388    fn native_capability_registry_aggregates_without_duplicating_root_peers() {
4389        let mut registry = NativeCapabilityRegistry::default();
4390        for (key, kind, id) in [
4391            ("devices:user-1", "devices", "user-1"),
4392            ("space:room-1", "space", "room-1"),
4393        ] {
4394            registry
4395                .register(key.to_string(), kind.to_string(), id.to_string(), false)
4396                .expect("capability registration");
4397            registry.registrations.get_mut(key).unwrap().desired_peers =
4398                vec![serde_json::json!({"deviceId": "shared-peer", "ticket": "ticket-a"})];
4399        }
4400
4401        let (first_revision, peers_json) = registry.aggregate_desired_peers().unwrap();
4402        let peers: Vec<serde_json::Value> = serde_json::from_str(&peers_json).unwrap();
4403        assert_eq!(first_revision, 1);
4404        assert_eq!(peers.len(), 1, "one physical peer must be dialed once");
4405
4406        registry
4407            .register(
4408                "devices:user-1".to_string(),
4409                "devices".to_string(),
4410                "user-1".to_string(),
4411                false,
4412            )
4413            .expect("exact registration is idempotent");
4414        assert!(registry
4415            .register(
4416                "devices:user-1".to_string(),
4417                "space".to_string(),
4418                "room-1".to_string(),
4419                false,
4420            )
4421            .is_err());
4422    }
4423
4424    #[test]
4425    fn native_capability_registry_merges_sparse_roles_deterministically() {
4426        fn aggregate(reverse_registration_order: bool) -> (u64, Vec<serde_json::Value>) {
4427            let mut registry = NativeCapabilityRegistry::default();
4428            let registrations = if reverse_registration_order {
4429                [
4430                    ("space:backup", "space", "backup"),
4431                    ("room:active", "room", "active"),
4432                ]
4433            } else {
4434                [
4435                    ("room:active", "room", "active"),
4436                    ("space:backup", "space", "backup"),
4437                ]
4438            };
4439            for (key, kind, id) in registrations {
4440                registry
4441                    .register(key.to_string(), kind.to_string(), id.to_string(), false)
4442                    .unwrap();
4443            }
4444            registry
4445                .registrations
4446                .get_mut("room:active")
4447                .unwrap()
4448                .desired_peers = vec![serde_json::json!({
4449                "deviceId": "shared-peer",
4450                "ticket": "active-ticket",
4451                "topologyRole": "active",
4452                "topologyRevision": 3,
4453            })];
4454            registry
4455                .registrations
4456                .get_mut("space:backup")
4457                .unwrap()
4458                .desired_peers = vec![serde_json::json!({
4459                "deviceId": "shared-peer",
4460                "ticket": "backup-ticket",
4461                "topologyRole": "backup",
4462                "topologyRevision": 100,
4463            })];
4464            let (revision, peers_json) = registry.aggregate_desired_peers().unwrap();
4465            (revision, serde_json::from_str(&peers_json).unwrap())
4466        }
4467
4468        let forward = aggregate(false);
4469        let reverse = aggregate(true);
4470        assert_eq!(
4471            forward, reverse,
4472            "registration order is not lifecycle authority"
4473        );
4474        assert_eq!(forward.0, 1);
4475        assert_eq!(forward.1.len(), 1);
4476        assert_eq!(forward.1[0]["ticket"], "active-ticket");
4477        assert_eq!(forward.1[0]["topologyRole"], "active");
4478        assert_eq!(forward.1[0]["topologyRevision"], 1);
4479    }
4480
4481    #[test]
4482    fn native_capability_registry_keeps_physical_peer_until_all_references_withdraw() {
4483        let mut registry = NativeCapabilityRegistry::default();
4484        for (key, kind) in [("room:a", "room"), ("space:b", "space")] {
4485            registry
4486                .register(key.to_string(), kind.to_string(), key.to_string(), false)
4487                .unwrap();
4488        }
4489        registry
4490            .registrations
4491            .get_mut("room:a")
4492            .unwrap()
4493            .desired_peers = vec![serde_json::json!({
4494            "deviceId": "shared-peer",
4495            "topologyRole": "backup",
4496            "topologyRevision": 41,
4497        })];
4498        registry
4499            .registrations
4500            .get_mut("space:b")
4501            .unwrap()
4502            .desired_peers = vec![serde_json::json!({
4503            "deviceId": "shared-peer",
4504            "nodeId": "shared-node",
4505            "ticket": "space-ticket",
4506            "topologyRole": "active",
4507            "topologyRevision": 7,
4508        })];
4509
4510        let (_, both_json) = registry.aggregate_desired_peers().unwrap();
4511        let both: Vec<serde_json::Value> = serde_json::from_str(&both_json).unwrap();
4512        assert_eq!(both.len(), 1);
4513        assert_eq!(both[0]["topologyRole"], "active");
4514        assert_eq!(both[0]["ticket"], "space-ticket");
4515        assert_eq!(
4516            registry.registrations["room:a"].desired_peers[0]["topologyRevision"], 41,
4517            "root projection must not rewrite avenue-local lease state",
4518        );
4519
4520        registry
4521            .registrations
4522            .get_mut("space:b")
4523            .unwrap()
4524            .desired_peers
4525            .clear();
4526        let (_, room_only_json) = registry.aggregate_desired_peers().unwrap();
4527        let room_only: Vec<serde_json::Value> = serde_json::from_str(&room_only_json).unwrap();
4528        assert_eq!(
4529            room_only.len(),
4530            1,
4531            "one remaining capability keeps the leg desired"
4532        );
4533        assert_eq!(room_only[0]["topologyRole"], "backup");
4534
4535        registry
4536            .registrations
4537            .get_mut("room:a")
4538            .unwrap()
4539            .desired_peers
4540            .clear();
4541        let (_, empty_json) = registry.aggregate_desired_peers().unwrap();
4542        let empty: Vec<serde_json::Value> = serde_json::from_str(&empty_json).unwrap();
4543        assert!(
4544            empty.is_empty(),
4545            "physical desire ends only after every reference withdraws"
4546        );
4547    }
4548
4549    #[test]
4550    fn native_capability_registry_scopes_peer_data_and_rejects_unscoped_streams() {
4551        let mut registry = NativeCapabilityRegistry::default();
4552        registry
4553            .register(
4554                "devices:user-1".to_string(),
4555                "devices".to_string(),
4556                "user-1".to_string(),
4557                false,
4558            )
4559            .unwrap();
4560        registry
4561            .registrations
4562            .get_mut("devices:user-1")
4563            .unwrap()
4564            .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
4565
4566        let explicit = openrtc::client::NativePeerDataEvent {
4567            connection_id: "peer-1".to_string(),
4568            remote_node_id: None,
4569            transport: "webrtc".to_string(),
4570            transport_stable_id: 1,
4571            transport_generation: 1,
4572            route_generation: 1,
4573            payload: serde_json::to_vec(&serde_json::json!({
4574                "capability": "devices:user-1",
4575                "kind": "raw",
4576                "payload": [1, 2, 3]
4577            }))
4578            .unwrap(),
4579        };
4580        assert_eq!(
4581            registry.capability_keys_for_peer_data(&explicit),
4582            vec!["devices:user-1".to_string()]
4583        );
4584
4585        let unknown = openrtc::client::NativePeerDataEvent {
4586            payload: serde_json::to_vec(&serde_json::json!({
4587                "capability": "space:unknown",
4588                "kind": "raw"
4589            }))
4590            .unwrap(),
4591            ..explicit.clone()
4592        };
4593        assert!(registry.capability_keys_for_peer_data(&unknown).is_empty());
4594
4595        registry
4596            .register(
4597                "space:shared".to_string(),
4598                "space".to_string(),
4599                "shared".to_string(),
4600                false,
4601            )
4602            .unwrap();
4603        registry
4604            .registrations
4605            .get_mut("space:shared")
4606            .unwrap()
4607            .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
4608        let unscoped = openrtc::client::NativePeerDataEvent {
4609            payload: serde_json::to_vec(&serde_json::json!({
4610                "kind": "raw",
4611                "payload": [1, 2, 3]
4612            }))
4613            .unwrap(),
4614            ..explicit.clone()
4615        };
4616        assert!(
4617            registry.capability_keys_for_peer_data(&unscoped).is_empty(),
4618            "a shared physical peer cannot fan unscoped data into two avenues",
4619        );
4620
4621        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
4622            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4623            metadata: None,
4624        };
4625        assert_eq!(
4626            registry.projected_stream_capability(&unscoped_channel),
4627            None,
4628            "stream ownership must never be inferred from capability count"
4629        );
4630    }
4631
4632    #[test]
4633    fn native_device_proof_uses_host_signer_without_exporting_private_key() {
4634        let state = OpenRtcTauriState::new(test_config())
4635            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4636        assert_eq!(
4637            state
4638                .native_device_key_signer
4639                .as_deref()
4640                .unwrap()
4641                .offline_assurance("app_native_test"),
4642            openrtc::offline::OfflineAssurance::Software
4643        );
4644        let public = state
4645            .device_public_key("app_native_test")
4646            .expect("public key");
4647        assert_eq!(public["kty"], "OKP");
4648        assert_eq!(public["crv"], "Ed25519");
4649        assert!(public.get("d").is_none());
4650        let signature = state
4651            .sign_device_proof("app_native_test", "openrtc:v2:test")
4652            .expect("signature");
4653        assert_eq!(
4654            base64::engine::general_purpose::URL_SAFE_NO_PAD
4655                .decode(signature)
4656                .expect("base64"),
4657            vec![7; 64]
4658        );
4659    }
4660
4661    #[test]
4662    fn offline_support_reports_signer_and_compiled_lan_truth() {
4663        let unavailable = OpenRtcTauriState::new(test_config()).offline_runtime_support();
4664        assert!(!unavailable.provisioning);
4665        assert_eq!(unavailable.local_mesh, cfg!(feature = "transport-lan"));
4666        assert_eq!(
4667            unavailable.reason,
4668            Some("native host device signer is unavailable")
4669        );
4670
4671        let available = OpenRtcTauriState::new(test_config())
4672            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()))
4673            .offline_runtime_support();
4674        assert!(available.provisioning);
4675        assert_eq!(available.local_mesh, cfg!(feature = "transport-lan"));
4676        assert_eq!(
4677            available.reason,
4678            (!cfg!(feature = "transport-lan"))
4679                .then_some("native host was built without transport-lan"),
4680        );
4681    }
4682
4683    #[tokio::test]
4684    async fn public_runtime_status_uses_host_facts_and_keeps_broadcast_unavailable() {
4685        use openrtc::client::CapabilityMaturity;
4686
4687        let unavailable_state = OpenRtcTauriState::new(test_config());
4688        let unavailable = unavailable_state
4689            .project_public_runtime_status(unavailable_state.client().runtime_status().await);
4690        assert_eq!(
4691            unavailable.product_maturity.offline_edge,
4692            CapabilityMaturity::Unavailable
4693        );
4694        assert_eq!(
4695            unavailable.product_maturity.broadcast,
4696            CapabilityMaturity::Unavailable
4697        );
4698        let serialized = serde_json::to_value(&unavailable).expect("serialized runtime status");
4699        assert_eq!(serialized["productMaturity"]["offlineEdge"], "unavailable");
4700        assert_eq!(serialized["productMaturity"]["broadcast"], "unavailable");
4701
4702        let installed_state = OpenRtcTauriState::new(test_config())
4703            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4704        let installed = installed_state
4705            .project_public_runtime_status(installed_state.client().runtime_status().await);
4706        assert_eq!(
4707            installed.product_maturity.offline_edge,
4708            if cfg!(feature = "transport-lan") {
4709                CapabilityMaturity::Preview
4710            } else {
4711                CapabilityMaturity::SupportOnly
4712            }
4713        );
4714        assert_eq!(
4715            installed.product_maturity.broadcast,
4716            CapabilityMaturity::Unavailable
4717        );
4718    }
4719
4720    #[test]
4721    fn native_certificate_records_round_trip_through_host_secure_storage() {
4722        let state = OpenRtcTauriState::new(test_config())
4723            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4724        let app_tag = "app_native_test";
4725        let key = "openrtc:v2:device-session:abc:device-1";
4726        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
4727        state
4728            .write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
4729            .unwrap();
4730        assert_eq!(
4731            state.read_secure_record(app_tag, key).unwrap().as_deref(),
4732            Some("{\"token\":\"bound\"}")
4733        );
4734        state.delete_secure_record(app_tag, key).unwrap();
4735        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
4736    }
4737
4738    #[test]
4739    fn native_device_proof_fails_closed_without_secure_host_signer() {
4740        let state = OpenRtcTauriState::new(test_config());
4741        assert!(state
4742            .device_public_key("app_native_test")
4743            .expect_err("missing signer must fail")
4744            .contains("secure-storage signer"));
4745    }
4746
4747    fn native_transport_install_context() -> InstallContext {
4748        InstallContext {
4749            data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
4750        }
4751    }
4752
4753    fn managed_session_test_result(device_id: &str) -> StartSessionResult {
4754        StartSessionResult {
4755            local_node_id: format!("node-{device_id}"),
4756            ticket_scope: Some("user-device".to_string()),
4757            ticket: None,
4758            presence_started: true,
4759            auto_connect_started: true,
4760            local_device: openrtc::native_device::NativeDeviceIdentity {
4761                device_id: device_id.to_string(),
4762                device_name: "Test Device".to_string(),
4763                created_at_ms: 1,
4764                updated_at_ms: 1,
4765                name_source: None,
4766                system_info: None,
4767            },
4768        }
4769    }
4770
4771    async fn simulate_managed_session_start(
4772        state: Arc<OpenRtcTauriState>,
4773        device_id: String,
4774        pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
4775        starts: Arc<AtomicUsize>,
4776        stopped_owners: Arc<Mutex<Vec<usize>>>,
4777    ) -> SessionDisposition {
4778        let _start_guard = state.managed_session_start_guard.lock().await;
4779        let client = state.client();
4780        let key = format!("{}:{device_id}", client.app_tag());
4781        let (disposition, active_epoch) = {
4782            let active = state.managed_session.lock().await;
4783            let disposition = managed_session_disposition(
4784                active.as_ref().map(|record| record.key.as_str()),
4785                active
4786                    .as_ref()
4787                    .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
4788                active
4789                    .as_ref()
4790                    .is_some_and(|record| record.result.presence_started),
4791                active
4792                    .as_ref()
4793                    .is_some_and(|record| record.result.auto_connect_started),
4794                &key,
4795                true,
4796                true,
4797            );
4798            (
4799                disposition,
4800                active.as_ref().map(|record| record.owner_epoch),
4801            )
4802        };
4803
4804        if disposition == SessionDisposition::Reuse {
4805            return disposition;
4806        }
4807
4808        if disposition == SessionDisposition::Replace {
4809            let previous = state
4810                .managed_session
4811                .lock()
4812                .await
4813                .take()
4814                .expect("replacement must have a previous owner");
4815            stopped_owners
4816                .lock()
4817                .await
4818                .push(Arc::as_ptr(&previous.owner_client) as usize);
4819        }
4820
4821        starts.fetch_add(1, Ordering::SeqCst);
4822        if let Some((entered, release)) = pause {
4823            entered.send(()).expect("start observer must be waiting");
4824            release.await.expect("start release must be sent");
4825        }
4826
4827        let owner_epoch = match disposition {
4828            SessionDisposition::Refresh => {
4829                active_epoch.expect("refresh must preserve the active owner epoch")
4830            }
4831            SessionDisposition::Start | SessionDisposition::Replace => {
4832                state.allocate_managed_session_owner_epoch()
4833            }
4834            SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
4835        };
4836        let result = managed_session_test_result(&device_id);
4837        *state.managed_session.lock().await = Some(ManagedSessionRecord {
4838            key,
4839            owner_client: client,
4840            owner_epoch,
4841            result,
4842        });
4843        disposition
4844    }
4845
4846    #[tokio::test]
4847    async fn concurrent_same_key_start_reuses_one_owner_epoch() {
4848        let state = Arc::new(OpenRtcTauriState::new(test_config()));
4849        let starts = Arc::new(AtomicUsize::new(0));
4850        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
4851        let (first_entered_tx, first_entered_rx) = oneshot::channel();
4852        let (first_release_tx, first_release_rx) = oneshot::channel();
4853
4854        let first = tokio::spawn(simulate_managed_session_start(
4855            state.clone(),
4856            "same-device".to_string(),
4857            Some((first_entered_tx, first_release_rx)),
4858            starts.clone(),
4859            stopped_owners.clone(),
4860        ));
4861        first_entered_rx.await.expect("first start must pause");
4862
4863        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
4864        let second_state = state.clone();
4865        let second_starts = starts.clone();
4866        let second_stopped_owners = stopped_owners.clone();
4867        let second = tokio::spawn(async move {
4868            second_attempting_tx.send(()).unwrap();
4869            simulate_managed_session_start(
4870                second_state,
4871                "same-device".to_string(),
4872                None,
4873                second_starts,
4874                second_stopped_owners,
4875            )
4876            .await
4877        });
4878        second_attempting_rx.await.unwrap();
4879        tokio::task::yield_now().await;
4880        assert_eq!(starts.load(Ordering::SeqCst), 1);
4881
4882        first_release_tx.send(()).unwrap();
4883        assert_eq!(
4884            first.await.expect("first start task"),
4885            SessionDisposition::Start
4886        );
4887        assert_eq!(
4888            second.await.expect("second start task"),
4889            SessionDisposition::Reuse
4890        );
4891        assert_eq!(starts.load(Ordering::SeqCst), 1);
4892
4893        let active = state
4894            .managed_session
4895            .lock()
4896            .await
4897            .clone()
4898            .expect("active owner");
4899        assert_eq!(active.owner_epoch, 1);
4900        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
4901        assert!(stopped_owners.lock().await.is_empty());
4902    }
4903
4904    #[tokio::test]
4905    async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
4906        let state = Arc::new(OpenRtcTauriState::new(test_config()));
4907        let starts = Arc::new(AtomicUsize::new(0));
4908        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
4909        let (first_entered_tx, first_entered_rx) = oneshot::channel();
4910        let (first_release_tx, first_release_rx) = oneshot::channel();
4911
4912        let first = tokio::spawn(simulate_managed_session_start(
4913            state.clone(),
4914            "first-device".to_string(),
4915            Some((first_entered_tx, first_release_rx)),
4916            starts.clone(),
4917            stopped_owners.clone(),
4918        ));
4919        first_entered_rx.await.expect("first start must pause");
4920        let first_owner = state.client();
4921
4922        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
4923        let second_state = state.clone();
4924        let second_starts = starts.clone();
4925        let second_stopped_owners = stopped_owners.clone();
4926        let second = tokio::spawn(async move {
4927            second_attempting_tx.send(()).unwrap();
4928            simulate_managed_session_start(
4929                second_state,
4930                "second-device".to_string(),
4931                None,
4932                second_starts,
4933                second_stopped_owners,
4934            )
4935            .await
4936        });
4937        second_attempting_rx.await.unwrap();
4938        tokio::task::yield_now().await;
4939        assert!(Arc::ptr_eq(&first_owner, &state.client()));
4940
4941        first_release_tx.send(()).unwrap();
4942        assert_eq!(
4943            first.await.expect("first start task"),
4944            SessionDisposition::Start
4945        );
4946        assert_eq!(
4947            second.await.expect("replacement start task"),
4948            SessionDisposition::Replace
4949        );
4950        assert_eq!(starts.load(Ordering::SeqCst), 2);
4951
4952        let active = state
4953            .managed_session
4954            .lock()
4955            .await
4956            .clone()
4957            .expect("active owner");
4958        assert_eq!(active.owner_epoch, 2);
4959        assert!(active.key.contains("second-device"));
4960        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
4961        assert_eq!(
4962            stopped_owners.lock().await.as_slice(),
4963            [Arc::as_ptr(&first_owner) as usize]
4964        );
4965        assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
4966    }
4967
4968    #[test]
4969    fn derives_app_tag_from_api_key_by_default() {
4970        let config = OpenRtcTauriConfig {
4971            api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
4972            ..OpenRtcTauriConfig::default()
4973        };
4974
4975        assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
4976    }
4977
4978    #[tokio::test]
4979    async fn requested_native_transport_without_installer_keeps_base_route_available() {
4980        let state = OpenRtcTauriState::new(test_config());
4981        let client = state.client();
4982        ensure_requested_native_transports(
4983            &state,
4984            &client,
4985            Some(&requested_ble_config()),
4986            &native_transport_install_context(),
4987        )
4988        .await
4989        .expect("an unavailable optional transport must not fail base Iroh startup");
4990
4991        assert!(state.installed_native_transports.lock().await.is_empty());
4992    }
4993
4994    #[tokio::test]
4995    async fn native_transport_installer_is_idempotent_for_one_client() {
4996        let installs = Arc::new(AtomicUsize::new(0));
4997        let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
4998            Arc::new(FakeTransportInstaller {
4999                installs: installs.clone(),
5000            }),
5001        );
5002        let client = state.client();
5003        let config = requested_ble_config();
5004
5005        ensure_requested_native_transports(
5006            &state,
5007            &client,
5008            Some(&config),
5009            &native_transport_install_context(),
5010        )
5011        .await
5012        .expect("first install");
5013        ensure_requested_native_transports(
5014            &state,
5015            &client,
5016            Some(&config),
5017            &native_transport_install_context(),
5018        )
5019        .await
5020        .expect("idempotent install");
5021
5022        assert_eq!(installs.load(Ordering::SeqCst), 1);
5023    }
5024
5025    #[tokio::test]
5026    async fn unavailable_native_transport_keeps_base_route_available() {
5027        let state = OpenRtcTauriState::new(test_config())
5028            .with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
5029        let client = state.client();
5030
5031        ensure_requested_native_transports(
5032            &state,
5033            &client,
5034            Some(&requested_ble_config()),
5035            &native_transport_install_context(),
5036        )
5037        .await
5038        .expect("optional transport installation failure must not fail base Iroh startup");
5039
5040        assert!(state.installed_native_transports.lock().await.is_empty());
5041    }
5042
5043    #[test]
5044    fn invalid_or_missing_api_key_fails_closed() {
5045        assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
5046        let mut config = OpenRtcTauriConfig::default();
5047        config.api_key = "pk_test_short".to_string();
5048        assert!(config.validated_api_key().is_err());
5049    }
5050
5051    #[cfg(feature = "native-broadcast-moq")]
5052    #[test]
5053    fn draft16_broadcast_feature_installs_one_private_native_adapter() {
5054        let state = OpenRtcTauriState::new(test_config());
5055        assert!(state.native_broadcast_adapter.is_some());
5056        assert!(state.native_broadcast_moq_adapter.is_some());
5057    }
5058
5059    #[test]
5060    fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
5061        let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
5062        let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
5063        let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
5064        let api_key = format!("pk_test_{}", "c".repeat(40));
5065        std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
5066        std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
5067        std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
5068
5069        let config = OpenRtcTauriConfig::from_env();
5070
5071        assert_eq!(config.api_key, api_key);
5072        assert_eq!(
5073            config.app_tag().unwrap(),
5074            openrtc::app_tag_from_api_key(&api_key)
5075        );
5076
5077        match previous_api_key {
5078            Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
5079            None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
5080        }
5081        match previous_project {
5082            Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
5083            None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
5084        }
5085        match previous_app_tag {
5086            Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
5087            None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
5088        }
5089    }
5090
5091    #[test]
5092    fn native_state_has_one_constructor_owned_app_identity() {
5093        let config = test_config();
5094        let expected = openrtc::app_tag_from_api_key(&config.api_key);
5095        let state = OpenRtcTauriState::new(config);
5096        let first = state.client();
5097        let second = state.client();
5098
5099        assert_eq!(first.app_tag(), expected);
5100        assert!(Arc::ptr_eq(&first, &second));
5101    }
5102
5103    #[test]
5104    fn managed_session_uses_persisted_native_device_id() {
5105        let identity = openrtc::native_device::NativeDeviceIdentity {
5106            device_id: "persisted-native-device".to_string(),
5107            device_name: "Mac".to_string(),
5108            created_at_ms: 1,
5109            updated_at_ms: 1,
5110            name_source: None,
5111            system_info: None,
5112        };
5113
5114        assert_eq!(
5115            managed_session_device_id(Some(" desktop-e2e-native "), &identity),
5116            "persisted-native-device"
5117        );
5118        assert_eq!(
5119            managed_session_device_id(None, &identity),
5120            "persisted-native-device"
5121        );
5122    }
5123
5124    #[test]
5125    fn repeated_managed_session_triggers_converge_on_one_owner() {
5126        let key = "app:persisted-native-device";
5127        let cases = [
5128            ("react-remount", true, true, true, true),
5129            ("hmr", true, true, true, true),
5130            ("auth-refresh", true, true, true, true),
5131            ("resume", true, true, true, true),
5132            ("alias-change", true, true, true, true),
5133        ];
5134        for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
5135            assert_eq!(
5136                managed_session_disposition(
5137                    Some(key),
5138                    true,
5139                    active_presence,
5140                    active_auto,
5141                    key,
5142                    requested_presence,
5143                    requested_auto,
5144                ),
5145                SessionDisposition::Reuse,
5146                "{label} must reuse the authoritative tuple"
5147            );
5148        }
5149
5150        assert_eq!(
5151            managed_session_disposition(Some(key), true, false, true, key, true, true),
5152            SessionDisposition::Refresh,
5153            "failed presence startup must retry idempotently"
5154        );
5155        assert_eq!(
5156            managed_session_disposition(
5157                Some(key),
5158                true,
5159                true,
5160                true,
5161                "app:other-native-device",
5162                true,
5163                true,
5164            ),
5165            SessionDisposition::Replace,
5166            "a physical native device change must replace the previous lifecycle owner"
5167        );
5168        assert_eq!(
5169            managed_session_disposition(Some(key), false, true, true, key, true, true),
5170            SessionDisposition::Replace,
5171            "the same tuple on a different client epoch must replace the previous owner"
5172        );
5173    }
5174
5175    #[test]
5176    fn user_device_revocation_invalidates_the_cached_managed_ticket() {
5177        assert!(revokes_managed_session("user-device"));
5178        assert!(revokes_managed_session("  user-device  "));
5179        assert!(!revokes_managed_session("share:example"));
5180        assert!(!revokes_managed_session(""));
5181    }
5182
5183    #[test]
5184    fn managed_session_metadata_uses_authoritative_device_id() {
5185        let raw = serde_json::json!({
5186            "deviceId": "persisted-native-device",
5187            "assistantDevice": {
5188                "deviceId": "desktop-e2e-native"
5189            }
5190        })
5191        .to_string();
5192
5193        let metadata =
5194            metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
5195        let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
5196
5197        assert_eq!(parsed["deviceId"], "desktop-e2e-native");
5198        assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
5199    }
5200}