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    if let Some(current) = subscription.as_ref() {
2945        if current.request_id == request_id {
2946            return Ok(request_id);
2947        }
2948        return Err(
2949            "OpenRTC native projection subscription already has a process owner".to_string(),
2950        );
2951    }
2952    *subscription = Some(NativeProjectionSubscription {
2953        request_id: request_id.clone(),
2954        channel,
2955    });
2956    Ok(request_id)
2957}
2958
2959#[tauri::command]
2960async fn get_projected_stream_diagnostics(
2961    state: tauri::State<'_, OpenRtcTauriState>,
2962) -> Result<ProjectedStreamDiagnostics, String> {
2963    Ok(ProjectedStreamDiagnostics {
2964        subscription_active: state.native_projection_subscription.lock().await.is_some(),
2965        offered: state.projected_stream_offered.load(Ordering::Relaxed),
2966        decoded: state.projected_stream_decoded.load(Ordering::Relaxed),
2967        projected: state.projected_stream_projected.load(Ordering::Relaxed),
2968        unhandled_non_channel: state
2969            .projected_stream_unhandled_non_channel
2970            .load(Ordering::Relaxed),
2971        unhandled_other_channel: state
2972            .projected_stream_unhandled_other_channel
2973            .load(Ordering::Relaxed),
2974        unhandled_no_subscription: state
2975            .projected_stream_unhandled_no_subscription
2976            .load(Ordering::Relaxed),
2977        failures: state.projected_stream_failures.load(Ordering::Relaxed),
2978        last_transport_stable_id: state
2979            .projected_stream_last_transport_stable_id
2980            .load(Ordering::Relaxed),
2981        validation_checks: state
2982            .projected_stream_validation_checks
2983            .load(Ordering::Relaxed),
2984        validation_rejections: state
2985            .projected_stream_validation_rejections
2986            .load(Ordering::Relaxed),
2987        last_validated_transport_stable_id: state
2988            .projected_stream_last_validated_transport_stable_id
2989            .load(Ordering::Relaxed),
2990        authorized: state.projected_stream_authorized.load(Ordering::Relaxed),
2991        unauthorized: state.projected_stream_unauthorized.load(Ordering::Relaxed),
2992        last_channel: state.projected_stream_last_channel.lock().await.clone(),
2993        last_protocol: state.projected_stream_last_protocol.lock().await.clone(),
2994        last_connection_id: state
2995            .projected_stream_last_connection_id
2996            .lock()
2997            .await
2998            .clone(),
2999        last_remote_node_id: state
3000            .projected_stream_last_remote_node_id
3001            .lock()
3002            .await
3003            .clone(),
3004    })
3005}
3006
3007#[tauri::command]
3008async fn is_current_transport_stable_id(
3009    state: tauri::State<'_, OpenRtcTauriState>,
3010    endpoint_id: String,
3011    transport_stable_id: u64,
3012) -> Result<bool, String> {
3013    state
3014        .projected_stream_validation_checks
3015        .fetch_add(1, Ordering::Relaxed);
3016    state
3017        .projected_stream_last_validated_transport_stable_id
3018        .store(transport_stable_id, Ordering::Relaxed);
3019    let is_current = state
3020        .client()
3021        .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
3022        .await
3023        .map_err(|error| error.to_string())?;
3024    if !is_current {
3025        state
3026            .projected_stream_validation_rejections
3027            .fetch_add(1, Ordering::Relaxed);
3028    }
3029    if native_stream_trace_enabled() {
3030        eprintln!(
3031            "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
3032            endpoint_id, transport_stable_id, is_current
3033        );
3034    }
3035    Ok(is_current)
3036}
3037
3038#[tauri::command]
3039async fn send_peer_message(
3040    state: tauri::State<'_, OpenRtcTauriState>,
3041    request: Request<'_>,
3042) -> Result<(), String> {
3043    let id = required_request_header(&request, PEER_ID_HEADER)?;
3044    let data = request_body_bytes(&request)?;
3045    state
3046        .client()
3047        .send_peer(&id, &data)
3048        .await
3049        .map_err(|error| error.to_string())
3050}
3051
3052#[tauri::command]
3053async fn encode_sparse_fanout_message(
3054    state: tauri::State<'_, OpenRtcTauriState>,
3055    request: Request<'_>,
3056) -> Result<Response, String> {
3057    let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
3058    let app_tag = required_request_header(&request, FANOUT_APP_TAG_HEADER)?;
3059    let payload = request_body_bytes(&request)?;
3060    let signing_request = state
3061        .client()
3062        .prepare_sparse_fanout_message(&capability, &payload)
3063        .await
3064        .map_err(|error| error.to_string())?;
3065    let signature = state.sign_device_message(&app_tag, &signing_request.signing_input)?;
3066    state
3067        .client()
3068        .finalize_sparse_fanout_message(&signing_request.request_id, &signature)
3069        .await
3070        .map(Response::new)
3071        .map_err(|error| error.to_string())
3072}
3073
3074#[tauri::command]
3075async fn accept_sparse_fanout_message(
3076    state: tauri::State<'_, OpenRtcTauriState>,
3077    request: Request<'_>,
3078) -> Result<serde_json::Value, String> {
3079    let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
3080    let source_peer = required_request_header(&request, FANOUT_SOURCE_PEER_HEADER)?;
3081    let encoded = request_body_bytes(&request)?;
3082    serde_json::to_value(
3083        state
3084            .client()
3085            .accept_sparse_fanout_message(&capability, &source_peer, &encoded)
3086            .await
3087            .map_err(|error| error.to_string())?,
3088    )
3089    .map_err(|error| error.to_string())
3090}
3091
3092#[tauri::command]
3093async fn sparse_fanout_diagnostics(
3094    state: tauri::State<'_, OpenRtcTauriState>,
3095    capability_key: String,
3096) -> Result<openrtc::sparse_fanout::SparseFanoutDiagnostics, String> {
3097    Ok(state
3098        .client()
3099        .sparse_fanout_diagnostics(&required_native_capability_part(
3100            capability_key,
3101            "capability key",
3102        )?)
3103        .await)
3104}
3105
3106#[derive(Debug, Serialize)]
3107#[serde(rename_all = "camelCase")]
3108struct NativeBroadcastPublisherDescriptor {
3109    handle: String,
3110    verification_key: Vec<u8>,
3111}
3112
3113#[derive(Debug, Serialize)]
3114#[serde(rename_all = "camelCase")]
3115struct NativeBroadcastSessionDescriptor {
3116    handle_id: String,
3117    id: String,
3118    role: openrtc::broadcast::BroadcastRole,
3119    state: openrtc::broadcast::BroadcastState,
3120    budget: openrtc::broadcast::BroadcastBudget,
3121}
3122
3123/// Creates a Rust-owned, sign-only publisher key. Only an opaque handle and
3124/// public verification key cross IPC; private bytes remain in this process.
3125#[tauri::command]
3126async fn prepare_openrtc_broadcast_publisher(
3127    state: tauri::State<'_, OpenRtcTauriState>,
3128) -> Result<NativeBroadcastPublisherDescriptor, String> {
3129    let signer = openrtc::broadcast::BroadcastPublisherSigner::generate()
3130        .map_err(|error| error.to_string())?;
3131    let verification_key = signer.verifying_key().to_vec();
3132    let handle = uuid::Uuid::new_v4().simple().to_string();
3133    state
3134        .broadcast_signers
3135        .lock()
3136        .await
3137        .insert(handle.clone(), signer);
3138    Ok(NativeBroadcastPublisherDescriptor {
3139        handle,
3140        verification_key,
3141    })
3142}
3143
3144#[tauri::command]
3145async fn release_openrtc_broadcast_publisher(
3146    state: tauri::State<'_, OpenRtcTauriState>,
3147    handle: String,
3148) -> Result<bool, String> {
3149    Ok(state
3150        .broadcast_signers
3151        .lock()
3152        .await
3153        .remove(&handle)
3154        .is_some())
3155}
3156
3157async fn drive_native_broadcast_actions(
3158    session: &openrtc::broadcast::BroadcastSession,
3159    adapter: &dyn NativeBroadcastAdapter,
3160) -> Result<(), String> {
3161    loop {
3162        let actions = session.take_actions(openrtc::broadcast::MAX_BROADCAST_ACTIONS);
3163        if actions.is_empty() {
3164            return Ok(());
3165        }
3166        for action in actions {
3167            if let Some(observation) = adapter.apply(session.clone(), action).await? {
3168                session
3169                    .observe(observation)
3170                    .map_err(|error| error.to_string())?;
3171            }
3172        }
3173    }
3174}
3175
3176async fn native_broadcast_driver(
3177    state: &OpenRtcTauriState,
3178    handle_id: &str,
3179) -> Result<Arc<Mutex<()>>, String> {
3180    state
3181        .broadcast_drivers
3182        .lock()
3183        .await
3184        .get(handle_id)
3185        .cloned()
3186        .ok_or_else(|| "broadcast session is unavailable".to_string())
3187}
3188
3189#[tauri::command]
3190async fn open_openrtc_broadcast(
3191    state: tauri::State<'_, OpenRtcTauriState>,
3192    grant_token: String,
3193    issuer_public_key: Vec<u8>,
3194    publisher_signer_handle: Option<String>,
3195    now_ms: u64,
3196) -> Result<NativeBroadcastSessionDescriptor, String> {
3197    let issuer_bytes: [u8; 32] = issuer_public_key
3198        .try_into()
3199        .map_err(|_| "broadcast issuer key must be 32 bytes".to_string())?;
3200    let issuer = VerifyingKey::from_bytes(&issuer_bytes)
3201        .map_err(|_| "broadcast issuer key is invalid".to_string())?;
3202    #[cfg(feature = "native-broadcast-moq")]
3203    let publication_verifying_key = if let Some(handle) = publisher_signer_handle.as_deref() {
3204        Some(
3205            state
3206                .broadcast_signers
3207                .lock()
3208                .await
3209                .get(handle)
3210                .ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?
3211                .verifying_key(),
3212        )
3213    } else {
3214        None
3215    };
3216    let broadcasts = state.client().broadcasts();
3217    let challenge = broadcasts
3218        .prepare_grant_verification(&grant_token, &issuer, now_ms)
3219        .map_err(|error| error.to_string())?;
3220    let device_signer = state.native_device_key_signer.as_deref().ok_or_else(|| {
3221        "native managed broadcast requires the host device-key signer".to_string()
3222    })?;
3223    let binding_signature = device_signer
3224        .sign(state.client().app_tag(), &challenge.signing_bytes)
3225        .map_err(|error| format!("sign broadcast installation challenge: {error}"))?;
3226    let grant = broadcasts
3227        .complete_grant_verification(
3228            &grant_token,
3229            &issuer,
3230            &challenge.handle,
3231            &binding_signature,
3232            now_ms,
3233        )
3234        .map_err(|error| error.to_string())?;
3235    let grant_generation = grant.grant_generation();
3236    #[cfg(feature = "native-broadcast-moq")]
3237    let native_moq_access = if state.native_broadcast_moq_adapter.is_some() {
3238        Some(
3239            state
3240                .resolve_native_broadcast_access(
3241                    &grant_token,
3242                    publication_verifying_key.as_ref(),
3243                    &issuer_bytes,
3244                )
3245                .await?,
3246        )
3247    } else {
3248        None
3249    };
3250    let session = if let Some(handle) = publisher_signer_handle.as_deref() {
3251        let signers = state.broadcast_signers.lock().await;
3252        let signer = signers
3253            .get(handle)
3254            .ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?;
3255        state
3256            .client()
3257            .broadcasts()
3258            .open_publisher(grant, signer, now_ms)
3259    } else {
3260        state.client().broadcasts().open(grant, now_ms)
3261    }
3262    .map_err(|error| error.to_string())?;
3263    let handle_id = format!("{}:{grant_generation}", session.id());
3264    #[cfg(feature = "native-broadcast-moq")]
3265    if let (Some(adapter), Some(access)) = (
3266        state.native_broadcast_moq_adapter.as_ref(),
3267        native_moq_access,
3268    ) {
3269        adapter
3270            .authorize(session.id(), grant_generation, access)
3271            .await;
3272    }
3273    let adapter = state.native_broadcast_adapter.as_deref().ok_or_else(|| {
3274        session.close();
3275        "native managed broadcast adapter is unavailable".to_string()
3276    })?;
3277    if let Err(error) = drive_native_broadcast_actions(&session, adapter).await {
3278        session.close();
3279        #[cfg(feature = "native-broadcast-moq")]
3280        if let Some(adapter) = state.native_broadcast_moq_adapter.as_ref() {
3281            adapter
3282                .forget_authorization(&session.id(), grant_generation)
3283                .await;
3284        }
3285        return Err(error);
3286    }
3287    let descriptor = NativeBroadcastSessionDescriptor {
3288        handle_id: handle_id.clone(),
3289        id: session.id(),
3290        role: session.role(),
3291        state: session.state(),
3292        budget: session.budget(),
3293    };
3294    state
3295        .broadcast_sessions
3296        .lock()
3297        .await
3298        .insert(handle_id.clone(), session);
3299    state
3300        .broadcast_drivers
3301        .lock()
3302        .await
3303        .insert(handle_id, Arc::new(Mutex::new(())));
3304    Ok(descriptor)
3305}
3306
3307#[tauri::command]
3308#[allow(clippy::too_many_arguments)]
3309async fn begin_openrtc_broadcast_publication(
3310    state: tauri::State<'_, OpenRtcTauriState>,
3311    handle_id: String,
3312    source_slot: String,
3313    publication_id: Option<String>,
3314    kind: String,
3315    codec: String,
3316    clock_rate: u32,
3317    coded_width: Option<u32>,
3318    coded_height: Option<u32>,
3319    channels: Option<u16>,
3320) -> Result<PortableMediaPublicationStart, String> {
3321    let driver = native_broadcast_driver(&state, &handle_id).await?;
3322    let _driver = driver.lock().await;
3323    let publication_id = match publication_id.as_deref().map(str::trim) {
3324        Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
3325        _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
3326    };
3327    let kind: openrtc::media::MediaKind = kind
3328        .parse()
3329        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3330    if !matches!(
3331        kind,
3332        openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
3333    ) {
3334        return Err("broadcast publications must be audio or video".to_string());
3335    }
3336    let publication = openrtc::media::MediaPublicationConfig {
3337        publication_id,
3338        media_generation: 0,
3339        kind,
3340        codec: codec
3341            .parse()
3342            .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?,
3343        clock_rate,
3344        coded_width,
3345        coded_height,
3346        channels,
3347    };
3348    let session = state
3349        .broadcast_sessions
3350        .lock()
3351        .await
3352        .get(&handle_id)
3353        .cloned()
3354        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3355    let publication = session
3356        .begin_publication(&source_slot, publication)
3357        .map_err(|error| error.to_string())?;
3358    let adapter = state
3359        .native_broadcast_adapter
3360        .as_deref()
3361        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3362    drive_native_broadcast_actions(&session, adapter).await?;
3363    Ok(PortableMediaPublicationStart {
3364        publication_id: publication.publication_id.to_string(),
3365        media_generation: publication.media_generation,
3366        // Broadcast control is signed and goes only to the private adapter.
3367        control: Vec::new(),
3368    })
3369}
3370
3371#[tauri::command]
3372#[allow(clippy::too_many_arguments)]
3373async fn publish_openrtc_broadcast_sample(
3374    state: tauri::State<'_, OpenRtcTauriState>,
3375    handle_id: String,
3376    publication_id: String,
3377    timestamp_us: u64,
3378    duration_us: u32,
3379    keyframe: bool,
3380    discardable: bool,
3381    payload: Vec<u8>,
3382) -> Result<(), String> {
3383    let driver = native_broadcast_driver(&state, &handle_id).await?;
3384    let _driver = driver.lock().await;
3385    openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
3386        .map_err(|error| error.to_string())?;
3387    let session = state
3388        .broadcast_sessions
3389        .lock()
3390        .await
3391        .get(&handle_id)
3392        .cloned()
3393        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3394    session
3395        .publish(
3396            parse_media_publication_id(&publication_id)?,
3397            openrtc::media::EncodedMediaSample {
3398                timestamp_us,
3399                duration_us,
3400                keyframe,
3401                discardable,
3402                payload,
3403            },
3404        )
3405        .map_err(|error| error.to_string())?;
3406    let adapter = state
3407        .native_broadcast_adapter
3408        .as_deref()
3409        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3410    drive_native_broadcast_actions(&session, adapter).await
3411}
3412
3413#[derive(Debug, Serialize)]
3414#[serde(rename_all = "camelCase")]
3415struct NativeBroadcastMediaSample {
3416    publication_id: String,
3417    media_generation: u32,
3418    sequence: u64,
3419    timestamp_us: u64,
3420    duration_us: u32,
3421    kind: &'static str,
3422    codec: &'static str,
3423    keyframe: bool,
3424    discardable: bool,
3425    payload: Vec<u8>,
3426}
3427
3428struct NativeBroadcastReceivedObject {
3429    source_slot: Option<String>,
3430    object: Vec<u8>,
3431}
3432
3433/// Retrieves only media objects already accepted by the private native
3434/// adapter and verified/decoded by the shared Rust broadcast owner. Provider
3435/// credentials, namespaces, and raw adapter actions never cross IPC.
3436#[tauri::command]
3437async fn receive_openrtc_broadcast_media(
3438    state: tauri::State<'_, OpenRtcTauriState>,
3439    handle_id: String,
3440    max: usize,
3441) -> Result<Vec<NativeBroadcastMediaSample>, String> {
3442    let driver = native_broadcast_driver(&state, &handle_id).await?;
3443    let _driver = driver.lock().await;
3444    let session = state
3445        .broadcast_sessions
3446        .lock()
3447        .await
3448        .get(&handle_id)
3449        .cloned()
3450        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3451    let adapter = state
3452        .native_broadcast_adapter
3453        .as_deref()
3454        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3455    drive_native_broadcast_actions(&session, adapter).await?;
3456    let objects: Vec<NativeBroadcastReceivedObject>;
3457    #[cfg(feature = "native-broadcast-moq")]
3458    {
3459        objects = if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
3460            native_moq
3461                .take_inbound_objects_with_source(&session, max.min(64))
3462                .await?
3463                .into_iter()
3464                .map(|item| NativeBroadcastReceivedObject {
3465                    source_slot: Some(item.source_slot),
3466                    object: item.object,
3467                })
3468                .collect()
3469        } else {
3470            adapter
3471                .take_inbound_objects(session.clone(), max.min(64))
3472                .await?
3473                .into_iter()
3474                .map(|object| NativeBroadcastReceivedObject {
3475                    source_slot: None,
3476                    object,
3477                })
3478                .collect()
3479        };
3480    }
3481    #[cfg(not(feature = "native-broadcast-moq"))]
3482    {
3483        objects = adapter
3484            .take_inbound_objects(session.clone(), max.min(64))
3485            .await?
3486            .into_iter()
3487            .map(|object| NativeBroadcastReceivedObject {
3488                source_slot: None,
3489                object,
3490            })
3491            .collect();
3492    }
3493    let mut accepted = Vec::with_capacity(objects.len());
3494    let mut delivered_bytes = 0_u64;
3495    for received in objects {
3496        let object_len = received.object.len() as u64;
3497        let chunk = if let Some(source_slot) = received.source_slot.as_deref() {
3498            let Some(media) = session
3499                .accept_media_for_source(source_slot, &received.object)
3500                .map_err(|error| error.to_string())?
3501            else {
3502                continue;
3503            };
3504            delivered_bytes = delivered_bytes.saturating_add(object_len);
3505            match media.media {
3506                Some(media) => Some(
3507                    openrtc::media::EncodedMediaChunk::decode(&media)
3508                        .map_err(|error| error.to_string())?,
3509                ),
3510                None => None,
3511            }
3512        } else {
3513            let chunk = session
3514                .accept_object(&received.object)
3515                .map_err(|error| error.to_string())?;
3516            delivered_bytes = delivered_bytes.saturating_add(object_len);
3517            chunk
3518        };
3519        let Some(chunk) = chunk else { continue };
3520        openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
3521            .and_then(|_| {
3522                openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
3523            })
3524            .map_err(|error| error.to_string())?;
3525        accepted.push(NativeBroadcastMediaSample {
3526            publication_id: chunk.publication_id.to_string(),
3527            media_generation: chunk.media_generation,
3528            sequence: chunk.sequence,
3529            timestamp_us: chunk.timestamp_us,
3530            duration_us: chunk.duration_us,
3531            kind: media_kind_name(chunk.kind),
3532            codec: media_codec_name(chunk.codec),
3533            keyframe: chunk.keyframe,
3534            discardable: chunk.discardable,
3535            payload: chunk.payload,
3536        });
3537    }
3538    #[cfg(feature = "native-broadcast-moq")]
3539    if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
3540        let usage = native_moq
3541            .record_delivered_bytes(&session, delivered_bytes)
3542            .await;
3543        let settlement = drive_native_broadcast_actions(&session, adapter).await;
3544        usage?;
3545        settlement?;
3546    }
3547    Ok(accepted)
3548}
3549
3550#[tauri::command]
3551async fn revoke_openrtc_broadcast(
3552    state: tauri::State<'_, OpenRtcTauriState>,
3553    handle_id: String,
3554    grant_generation: u64,
3555) -> Result<(), String> {
3556    let driver = native_broadcast_driver(&state, &handle_id).await?;
3557    let _driver = driver.lock().await;
3558    let session = state
3559        .broadcast_sessions
3560        .lock()
3561        .await
3562        .get(&handle_id)
3563        .cloned()
3564        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3565    session
3566        .revoke(grant_generation)
3567        .map_err(|error| error.to_string())?;
3568    let adapter = state
3569        .native_broadcast_adapter
3570        .as_deref()
3571        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3572    drive_native_broadcast_actions(&session, adapter).await
3573}
3574
3575#[tauri::command]
3576async fn close_openrtc_broadcast(
3577    state: tauri::State<'_, OpenRtcTauriState>,
3578    handle_id: String,
3579) -> Result<(), String> {
3580    let driver = native_broadcast_driver(&state, &handle_id).await?;
3581    let _driver = driver.lock().await;
3582    let session = state
3583        .broadcast_sessions
3584        .lock()
3585        .await
3586        .remove(&handle_id)
3587        .ok_or_else(|| "broadcast session is unavailable".to_string())?;
3588    session.close();
3589    let adapter = state
3590        .native_broadcast_adapter
3591        .as_deref()
3592        .ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
3593    let result = drive_native_broadcast_actions(&session, adapter).await;
3594    state.broadcast_drivers.lock().await.remove(&handle_id);
3595    result
3596}
3597
3598#[tauri::command]
3599async fn record_sparse_fanout_forward_queue_drop(
3600    state: tauri::State<'_, OpenRtcTauriState>,
3601    capability_key: String,
3602    count: u64,
3603) -> Result<(), String> {
3604    state
3605        .client()
3606        .record_sparse_fanout_forward_queue_drop(
3607            &required_native_capability_part(capability_key, "capability key")?,
3608            count,
3609        )
3610        .await
3611        .map_err(|error| error.to_string())
3612}
3613
3614#[tauri::command]
3615async fn is_peer_connected(
3616    state: tauri::State<'_, OpenRtcTauriState>,
3617    node_id: String,
3618) -> Result<bool, String> {
3619    state
3620        .client()
3621        .is_connected_str(&node_id)
3622        .await
3623        .map_err(|error| error.to_string())
3624}
3625
3626#[derive(Debug, Serialize)]
3627#[serde(rename_all = "camelCase")]
3628struct PortableMediaPublicationStart {
3629    publication_id: String,
3630    media_generation: u32,
3631    control: Vec<u8>,
3632}
3633
3634fn media_kind_name(kind: openrtc::media::MediaKind) -> &'static str {
3635    match kind {
3636        openrtc::media::MediaKind::Audio => "audio",
3637        openrtc::media::MediaKind::Video => "video",
3638        openrtc::media::MediaKind::Screen => "screen",
3639        openrtc::media::MediaKind::Data => "data",
3640    }
3641}
3642
3643fn media_codec_name(codec: openrtc::media::MediaCodec) -> &'static str {
3644    match codec {
3645        openrtc::media::MediaCodec::Opus => "opus",
3646        openrtc::media::MediaCodec::H264 => "h264",
3647        openrtc::media::MediaCodec::Vp8 => "vp8",
3648        openrtc::media::MediaCodec::Vp9 => "vp9",
3649        openrtc::media::MediaCodec::Av1 => "av1",
3650        openrtc::media::MediaCodec::Pcm => "pcm",
3651        openrtc::media::MediaCodec::Opaque => "opaque",
3652    }
3653}
3654
3655fn parse_media_publication_id(value: &str) -> Result<openrtc::media::PublicationId, String> {
3656    value
3657        .trim()
3658        .parse()
3659        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())
3660}
3661
3662#[tauri::command]
3663#[allow(clippy::too_many_arguments)]
3664async fn begin_openrtc_media_publication(
3665    state: tauri::State<'_, OpenRtcTauriState>,
3666    publication_id: Option<String>,
3667    kind: String,
3668    codec: String,
3669    clock_rate: u32,
3670    coded_width: Option<u32>,
3671    coded_height: Option<u32>,
3672    channels: Option<u16>,
3673) -> Result<PortableMediaPublicationStart, String> {
3674    let publication_id = match publication_id.as_deref().map(str::trim) {
3675        Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
3676        _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
3677    };
3678    let kind: openrtc::media::MediaKind = kind
3679        .parse()
3680        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3681    if !matches!(
3682        kind,
3683        openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
3684    ) {
3685        return Err("browser media publications must be audio or video".to_string());
3686    }
3687    let codec = codec
3688        .parse()
3689        .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
3690    let publication = openrtc::media::MediaPublicationConfig {
3691        publication_id,
3692        media_generation: 0,
3693        kind,
3694        codec,
3695        clock_rate,
3696        coded_width,
3697        coded_height,
3698        channels,
3699    };
3700    let (publication, control) = state
3701        .portable_media
3702        .lock()
3703        .await
3704        .begin_publication(publication)
3705        .map_err(|error| error.to_string())?;
3706    Ok(PortableMediaPublicationStart {
3707        publication_id: publication.publication_id.to_string(),
3708        media_generation: publication.media_generation,
3709        control,
3710    })
3711}
3712
3713#[tauri::command]
3714#[allow(clippy::too_many_arguments)]
3715async fn encode_openrtc_media_sample(
3716    state: tauri::State<'_, OpenRtcTauriState>,
3717    request: Request<'_>,
3718) -> Result<Response, String> {
3719    let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
3720    let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
3721    openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
3722        .map_err(|error| error.to_string())?;
3723    let duration_us = parsed_request_header(&request, MEDIA_DURATION_US_HEADER)?;
3724    let keyframe = parsed_request_header(&request, MEDIA_KEYFRAME_HEADER)?;
3725    let discardable = parsed_request_header(&request, MEDIA_DISCARDABLE_HEADER)?;
3726    let payload = request_body_bytes(&request)?;
3727    let encoded = state
3728        .portable_media
3729        .lock()
3730        .await
3731        .encode_sample(
3732            parse_media_publication_id(&publication_id)?,
3733            openrtc::media::EncodedMediaSample {
3734                timestamp_us,
3735                duration_us,
3736                keyframe,
3737                discardable,
3738                payload,
3739            },
3740        )
3741        .map_err(|error| error.to_string())?;
3742    Ok(Response::new(encoded))
3743}
3744
3745#[tauri::command]
3746async fn pause_openrtc_media_publication(
3747    state: tauri::State<'_, OpenRtcTauriState>,
3748    publication_id: String,
3749) -> Result<(), String> {
3750    state
3751        .portable_media
3752        .lock()
3753        .await
3754        .pause_publication(parse_media_publication_id(&publication_id)?)
3755        .map_err(|error| error.to_string())
3756}
3757
3758#[tauri::command]
3759async fn retire_openrtc_media_publication(
3760    state: tauri::State<'_, OpenRtcTauriState>,
3761    publication_id: String,
3762) -> Result<(), String> {
3763    state
3764        .portable_media
3765        .lock()
3766        .await
3767        .retire_publication(parse_media_publication_id(&publication_id)?);
3768    Ok(())
3769}
3770
3771#[tauri::command]
3772async fn retire_openrtc_media_receiver(
3773    state: tauri::State<'_, OpenRtcTauriState>,
3774    publication_id: String,
3775    media_generation: u32,
3776) -> Result<(), String> {
3777    state.portable_media.lock().await.retire_receiver(
3778        parse_media_publication_id(&publication_id)?,
3779        media_generation,
3780    );
3781    Ok(())
3782}
3783
3784#[tauri::command]
3785async fn decode_openrtc_media_chunk(
3786    state: tauri::State<'_, OpenRtcTauriState>,
3787    request: Request<'_>,
3788) -> Result<Response, String> {
3789    let encoded = request_body_bytes(&request)?;
3790    let chunk = state
3791        .portable_media
3792        .lock()
3793        .await
3794        .decode_chunk(&encoded)
3795        .map_err(|error| error.to_string())?;
3796    openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
3797        .and_then(|_| {
3798            openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
3799        })
3800        .map_err(|error| error.to_string())?;
3801    let metadata = serde_json::to_vec(&serde_json::json!({
3802        "publicationId": chunk.publication_id.to_string(),
3803        "mediaGeneration": chunk.media_generation,
3804        "sequence": chunk.sequence,
3805        "timestampUs": chunk.timestamp_us,
3806        "durationUs": chunk.duration_us,
3807        "kind": media_kind_name(chunk.kind),
3808        "codec": media_codec_name(chunk.codec),
3809        "keyframe": chunk.keyframe,
3810        "discardable": chunk.discardable,
3811    }))
3812    .map_err(|error| format!("encode media IPC metadata failed: {error}"))?;
3813    let metadata_len =
3814        u32::try_from(metadata.len()).map_err(|_| "media IPC metadata is too large".to_string())?;
3815    let mut response = Vec::with_capacity(4 + metadata.len() + chunk.payload.len());
3816    response.extend_from_slice(&metadata_len.to_be_bytes());
3817    response.extend_from_slice(&metadata);
3818    response.extend_from_slice(&chunk.payload);
3819    Ok(Response::new(response))
3820}
3821
3822#[tauri::command]
3823async fn decode_openrtc_media_control(
3824    state: tauri::State<'_, OpenRtcTauriState>,
3825    request: Request<'_>,
3826) -> Result<serde_json::Value, String> {
3827    let encoded = request_body_bytes(&request)?;
3828    let control = state
3829        .portable_media
3830        .lock()
3831        .await
3832        .decode_control(&encoded)
3833        .map_err(|error| error.to_string())?;
3834    Ok(match control {
3835        openrtc::media::MediaControlFrame::Publish {
3836            publication_id,
3837            media_generation,
3838            kind,
3839            codec,
3840            clock_rate,
3841            coded_width,
3842            coded_height,
3843            channels,
3844        } => serde_json::json!({
3845            "type": "publish",
3846            "publicationId": publication_id.to_string(),
3847            "mediaGeneration": media_generation,
3848            "kind": media_kind_name(kind),
3849            "codec": media_codec_name(codec),
3850            "clockRate": clock_rate,
3851            "codedWidth": coded_width,
3852            "codedHeight": coded_height,
3853            "channels": channels,
3854        }),
3855        openrtc::media::MediaControlFrame::SetEnabled {
3856            publication_id,
3857            media_generation,
3858            enabled,
3859        } => serde_json::json!({
3860            "type": "set-enabled",
3861            "publicationId": publication_id.to_string(),
3862            "mediaGeneration": media_generation,
3863            "enabled": enabled,
3864        }),
3865        openrtc::media::MediaControlFrame::RequestKeyframe {
3866            publication_id,
3867            media_generation,
3868        } => serde_json::json!({
3869            "type": "request-keyframe",
3870            "publicationId": publication_id.to_string(),
3871            "mediaGeneration": media_generation,
3872        }),
3873        openrtc::media::MediaControlFrame::Stop {
3874            publication_id,
3875            media_generation,
3876            reason,
3877        } => serde_json::json!({
3878            "type": "stop",
3879            "publicationId": publication_id.to_string(),
3880            "mediaGeneration": media_generation,
3881            "reason": reason,
3882        }),
3883    })
3884}
3885
3886#[tauri::command]
3887async fn open_peer_bi_stream<R: Runtime>(
3888    app: tauri::AppHandle<R>,
3889    state: tauri::State<'_, OpenRtcTauriState>,
3890    peer_id: String,
3891    timeout_ms: Option<u64>,
3892) -> Result<OpenBiResult, String> {
3893    open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
3894        client.open_peer_bi(&peer_id, timeout_ms).await
3895    })
3896    .await
3897}
3898
3899#[tauri::command]
3900async fn open_peer_bi_transport_only_stream<R: Runtime>(
3901    app: tauri::AppHandle<R>,
3902    state: tauri::State<'_, OpenRtcTauriState>,
3903    peer_id: String,
3904    timeout_ms: Option<u64>,
3905) -> Result<OpenBiResult, String> {
3906    open_peer_bi_with(
3907        app,
3908        state,
3909        peer_id,
3910        timeout_ms,
3911        |client, peer_id, timeout_ms| async move {
3912            client
3913                .open_peer_bi_transport_only(&peer_id, timeout_ms)
3914                .await
3915                .map(|(connection_id, remote_node_id, send, recv)| {
3916                    (
3917                        connection_id,
3918                        remote_node_id,
3919                        PeerSendStream::plain(send),
3920                        PeerRecvStream::plain(recv),
3921                    )
3922                })
3923        },
3924    )
3925    .await
3926}
3927
3928#[tauri::command]
3929async fn open_peer_uni_stream(
3930    state: tauri::State<'_, OpenRtcTauriState>,
3931    peer_id: String,
3932    timeout_ms: Option<u64>,
3933) -> Result<OpenUniResult, String> {
3934    let peer_id = peer_id.trim().to_string();
3935    if peer_id.is_empty() {
3936        return Err("peerId is required".to_string());
3937    }
3938    let (connection_id, remote_node_id, send) = state
3939        .client()
3940        .open_peer_uni(&peer_id, timeout_ms)
3941        .await
3942        .map_err(|error| format!("open peer uni stream failed: {error}"))?;
3943    let stream_id = uuid::Uuid::new_v4().to_string();
3944    state
3945        .peer_uni_streams
3946        .lock()
3947        .await
3948        .insert(stream_id.clone(), Arc::new(Mutex::new(Some(send))));
3949    Ok(OpenUniResult {
3950        stream_id,
3951        connection_id,
3952        remote_node_id,
3953    })
3954}
3955
3956#[tauri::command]
3957async fn write_peer_bi_stream(
3958    state: tauri::State<'_, OpenRtcTauriState>,
3959    request: Request<'_>,
3960) -> Result<(), String> {
3961    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
3962    let bytes = request_body_bytes(&request)?;
3963
3964    let send = {
3965        let streams = state.peer_bi_streams.lock().await;
3966        streams
3967            .get(&stream_id)
3968            .map(|handle| handle.send.clone())
3969            .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
3970    };
3971    let result = send
3972        .lock()
3973        .await
3974        .as_mut()
3975        .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
3976        .write_all(&bytes)
3977        .await
3978        .map_err(|error| format!("write peer bi stream failed: {error}"));
3979    result
3980}
3981
3982#[tauri::command]
3983async fn start_peer_bi_stream_read<R: Runtime>(
3984    _app: tauri::AppHandle<R>,
3985    state: tauri::State<'_, OpenRtcTauriState>,
3986    stream_id: String,
3987    channel: Channel<Response>,
3988) -> Result<(), String> {
3989    let stream_id = stream_id.trim().to_string();
3990    if stream_id.is_empty() {
3991        return Err("streamId is required".to_string());
3992    }
3993
3994    let mut streams = state.peer_bi_streams.lock().await;
3995    let handle = streams
3996        .get_mut(&stream_id)
3997        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
3998    if handle.read_task.is_some() {
3999        return Ok(());
4000    }
4001    let recv = handle
4002        .recv
4003        .take()
4004        .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
4005    handle.read_task = Some(spawn_peer_bi_stream_reader(channel, stream_id, recv));
4006    Ok(())
4007}
4008
4009/// Cancel only the receive half of a bidirectional stream.
4010///
4011/// Web Streams permits a consumer to cancel an unused readable while keeping
4012/// its writable alive. Media publications rely on that half-close behavior, so
4013/// the native IPC projection must not remove or finish the Rust send handle.
4014#[tauri::command]
4015async fn cancel_peer_bi_stream_read(
4016    state: tauri::State<'_, OpenRtcTauriState>,
4017    stream_id: String,
4018) -> Result<(), String> {
4019    let stream_id = stream_id.trim().to_string();
4020    if stream_id.is_empty() {
4021        return Err("streamId is required".to_string());
4022    }
4023
4024    let mut streams = state.peer_bi_streams.lock().await;
4025    let handle = streams
4026        .get_mut(&stream_id)
4027        .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
4028    handle.recv.take();
4029    if let Some(read_task) = handle.read_task.take() {
4030        read_task.abort();
4031    }
4032    Ok(())
4033}
4034
4035#[tauri::command]
4036async fn close_peer_bi_stream(
4037    state: tauri::State<'_, OpenRtcTauriState>,
4038    stream_id: String,
4039) -> Result<(), String> {
4040    let stream_id = stream_id.trim().to_string();
4041    if stream_id.is_empty() {
4042        return Err("streamId is required".to_string());
4043    }
4044
4045    let handle = {
4046        let mut streams = state.peer_bi_streams.lock().await;
4047        streams.remove(&stream_id)
4048    };
4049    if let Some(handle) = handle {
4050        // Hold the writer mutex before finishing so close waits for any
4051        // in-flight chunk write instead of skipping the flush when Arc clones exist.
4052        let mut send = handle.send.lock().await;
4053        if let Some(send) = send.take() {
4054            let _ = send.finish();
4055        }
4056        drop(send);
4057        if let Some(read_task) = handle.read_task {
4058            read_task.abort();
4059        }
4060    }
4061    Ok(())
4062}
4063
4064#[tauri::command]
4065async fn write_peer_uni_stream(
4066    state: tauri::State<'_, OpenRtcTauriState>,
4067    request: Request<'_>,
4068) -> Result<(), String> {
4069    let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
4070    let bytes = request_body_bytes(&request)?;
4071    let send = {
4072        let streams = state.peer_uni_streams.lock().await;
4073        streams
4074            .get(&stream_id)
4075            .cloned()
4076            .ok_or_else(|| format!("peer uni stream not found: {stream_id}"))?
4077    };
4078    let result = send
4079        .lock()
4080        .await
4081        .as_mut()
4082        .ok_or_else(|| format!("peer uni stream already closed: {stream_id}"))?
4083        .write_all(&bytes)
4084        .await
4085        .map_err(|error| format!("write peer uni stream failed: {error}"));
4086    result
4087}
4088
4089#[tauri::command]
4090async fn close_peer_uni_stream(
4091    state: tauri::State<'_, OpenRtcTauriState>,
4092    stream_id: String,
4093) -> Result<(), String> {
4094    let stream_id = stream_id.trim().to_string();
4095    if stream_id.is_empty() {
4096        return Err("streamId is required".to_string());
4097    }
4098    let send = state.peer_uni_streams.lock().await.remove(&stream_id);
4099    if let Some(send) = send {
4100        if let Some(send) = send.lock().await.take() {
4101            let _ = send.finish();
4102        }
4103    }
4104    Ok(())
4105}
4106
4107#[cfg(test)]
4108mod tests {
4109    use super::*;
4110    use std::collections::HashMap as StdHashMap;
4111    use std::sync::atomic::{AtomicUsize, Ordering};
4112    use std::sync::Mutex as StdMutex;
4113    use tokio::sync::oneshot;
4114
4115    struct FakeTransportInstaller {
4116        installs: Arc<AtomicUsize>,
4117    }
4118
4119    #[derive(Default)]
4120    struct FakeDeviceKeySigner {
4121        records: StdMutex<StdHashMap<String, String>>,
4122    }
4123
4124    impl DeviceKeySigner for FakeDeviceKeySigner {
4125        fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
4126            Ok(serde_json::json!({
4127                "kty": "OKP",
4128                "crv": "Ed25519",
4129                "x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
4130            }))
4131        }
4132
4133        fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
4134            assert_eq!(challenge, b"openrtc:v2:test");
4135            Ok(vec![7; 64])
4136        }
4137
4138        fn delete(&self, _app_tag: &str) -> Result<(), String> {
4139            Ok(())
4140        }
4141
4142        fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
4143            Ok(self
4144                .records
4145                .lock()
4146                .unwrap()
4147                .get(&format!("{app_tag}:{key}"))
4148                .cloned())
4149        }
4150
4151        fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
4152            self.records
4153                .lock()
4154                .unwrap()
4155                .insert(format!("{app_tag}:{key}"), value.to_string());
4156            Ok(())
4157        }
4158
4159        fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
4160            self.records
4161                .lock()
4162                .unwrap()
4163                .remove(&format!("{app_tag}:{key}"));
4164            Ok(())
4165        }
4166    }
4167
4168    impl TransportInstaller for FakeTransportInstaller {
4169        fn id(&self) -> &'static str {
4170            "ble"
4171        }
4172
4173        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
4174            config.ble.as_ref().is_some_and(|ble| ble.enabled)
4175        }
4176
4177        fn install(
4178            &self,
4179            _client: Arc<openrtc::client::Client>,
4180            _config: openrtc::client::TransportConfig,
4181            _context: InstallContext,
4182        ) -> InstallFuture {
4183            self.installs.fetch_add(1, Ordering::SeqCst);
4184            Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
4185        }
4186    }
4187
4188    #[derive(Debug)]
4189    struct UnavailableTransportInstaller;
4190
4191    impl TransportInstaller for UnavailableTransportInstaller {
4192        fn id(&self) -> &'static str {
4193            "ble"
4194        }
4195
4196        fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
4197            config.ble.as_ref().is_some_and(|ble| ble.enabled)
4198        }
4199
4200        fn install(
4201            &self,
4202            _client: Arc<openrtc::client::Client>,
4203            _config: openrtc::client::TransportConfig,
4204            _context: InstallContext,
4205        ) -> InstallFuture {
4206            Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
4207        }
4208    }
4209
4210    fn requested_ble_config() -> openrtc::client::TransportConfig {
4211        openrtc::client::TransportConfig {
4212            ble: Some(openrtc::client::BleConfig {
4213                enabled: true,
4214                ..openrtc::client::BleConfig::default()
4215            }),
4216            ..openrtc::client::TransportConfig::default()
4217        }
4218    }
4219
4220    fn test_config() -> OpenRtcTauriConfig {
4221        OpenRtcTauriConfig {
4222            api_key: format!("pk_test_{}", "a".repeat(40)),
4223            ..OpenRtcTauriConfig::default()
4224        }
4225    }
4226
4227    #[test]
4228    fn native_room_architecture_requires_fixed_intent_until_rust_owns_the_gateway_handle() {
4229        let mut registry = NativeCapabilityRegistry::default();
4230        registry
4231            .register_with_architecture(
4232                "room:adaptive".to_string(),
4233                "room".to_string(),
4234                "adaptive".to_string(),
4235                true,
4236                Some(openrtc::native::RoomArchitectureMode::Auto),
4237            )
4238            .expect("room registration");
4239        assert!(
4240            !registry.registrations["room:adaptive"].uses_sparse_fanout(),
4241            "caller-supplied auto state cannot promote sparse fanout",
4242        );
4243        registry
4244            .register_with_architecture(
4245                "room:fixed-sparse".to_string(),
4246                "room".to_string(),
4247                "fixed-sparse".to_string(),
4248                false,
4249                Some(openrtc::native::RoomArchitectureMode::Sparse),
4250            )
4251            .expect("fixed sparse room registration");
4252        assert!(registry.registrations["room:fixed-sparse"].uses_sparse_fanout());
4253    }
4254
4255    #[test]
4256    fn native_room_architecture_rejects_non_room_and_respects_fixed_mesh_intent() {
4257        let mut registry = NativeCapabilityRegistry::default();
4258        assert!(registry
4259            .register_with_architecture(
4260                "space:not-room".to_string(),
4261                "space".to_string(),
4262                "not-room".to_string(),
4263                false,
4264                Some(openrtc::native::RoomArchitectureMode::Auto),
4265            )
4266            .is_err());
4267        registry
4268            .register_with_architecture(
4269                "room:fixed".to_string(),
4270                "room".to_string(),
4271                "fixed".to_string(),
4272                false,
4273                Some(openrtc::native::RoomArchitectureMode::Mesh),
4274            )
4275            .expect("fixed room registration");
4276        assert!(!registry.registrations["room:fixed"].uses_sparse_fanout());
4277    }
4278
4279    #[test]
4280    fn native_capability_registry_isolates_duplicate_peer_channels() {
4281        let mut registry = NativeCapabilityRegistry::default();
4282        registry
4283            .register(
4284                "devices:user-1".to_string(),
4285                "devices".to_string(),
4286                "user-1".to_string(),
4287                false,
4288            )
4289            .expect("devices registration");
4290        registry
4291            .register(
4292                "space:room-1".to_string(),
4293                "space".to_string(),
4294                "room-1".to_string(),
4295                false,
4296            )
4297            .expect("space registration");
4298        registry
4299            .registrations
4300            .get_mut("devices:user-1")
4301            .unwrap()
4302            .desired_peers = vec![serde_json::json!({
4303            "deviceId": "shared-peer",
4304            "nodeId": "device-route"
4305        })];
4306        registry
4307            .registrations
4308            .get_mut("space:room-1")
4309            .unwrap()
4310            .desired_peers = vec![serde_json::json!({
4311            "deviceId": "shared-peer",
4312            "nodeId": "space-route"
4313        })];
4314
4315        assert_eq!(
4316            registry.capability_keys_for_identity(
4317                Some("shared-peer"),
4318                Some("shared-peer"),
4319                None,
4320                None,
4321            ),
4322            vec!["devices:user-1".to_string(), "space:room-1".to_string()]
4323        );
4324
4325        let explicit_channel = openrtc::stream_metadata::ChannelMetadata {
4326            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4327            metadata: Some(serde_json::Map::from_iter([(
4328                "openrtcCapability".to_string(),
4329                serde_json::json!("space:room-1"),
4330            )])),
4331        };
4332        assert_eq!(
4333            registry.projected_stream_capability(&explicit_channel),
4334            Some("space:room-1".to_string())
4335        );
4336
4337        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
4338            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4339            metadata: None,
4340        };
4341        assert_eq!(
4342            registry.projected_stream_capability(&unscoped_channel),
4343            None
4344        );
4345    }
4346
4347    #[test]
4348    fn native_capability_registry_aggregates_without_duplicating_root_peers() {
4349        let mut registry = NativeCapabilityRegistry::default();
4350        for (key, kind, id) in [
4351            ("devices:user-1", "devices", "user-1"),
4352            ("space:room-1", "space", "room-1"),
4353        ] {
4354            registry
4355                .register(key.to_string(), kind.to_string(), id.to_string(), false)
4356                .expect("capability registration");
4357            registry.registrations.get_mut(key).unwrap().desired_peers =
4358                vec![serde_json::json!({"deviceId": "shared-peer", "ticket": "ticket-a"})];
4359        }
4360
4361        let (first_revision, peers_json) = registry.aggregate_desired_peers().unwrap();
4362        let peers: Vec<serde_json::Value> = serde_json::from_str(&peers_json).unwrap();
4363        assert_eq!(first_revision, 1);
4364        assert_eq!(peers.len(), 1, "one physical peer must be dialed once");
4365
4366        registry
4367            .register(
4368                "devices:user-1".to_string(),
4369                "devices".to_string(),
4370                "user-1".to_string(),
4371                false,
4372            )
4373            .expect("exact registration is idempotent");
4374        assert!(registry
4375            .register(
4376                "devices:user-1".to_string(),
4377                "space".to_string(),
4378                "room-1".to_string(),
4379                false,
4380            )
4381            .is_err());
4382    }
4383
4384    #[test]
4385    fn native_capability_registry_merges_sparse_roles_deterministically() {
4386        fn aggregate(reverse_registration_order: bool) -> (u64, Vec<serde_json::Value>) {
4387            let mut registry = NativeCapabilityRegistry::default();
4388            let registrations = if reverse_registration_order {
4389                [
4390                    ("space:backup", "space", "backup"),
4391                    ("room:active", "room", "active"),
4392                ]
4393            } else {
4394                [
4395                    ("room:active", "room", "active"),
4396                    ("space:backup", "space", "backup"),
4397                ]
4398            };
4399            for (key, kind, id) in registrations {
4400                registry
4401                    .register(key.to_string(), kind.to_string(), id.to_string(), false)
4402                    .unwrap();
4403            }
4404            registry
4405                .registrations
4406                .get_mut("room:active")
4407                .unwrap()
4408                .desired_peers = vec![serde_json::json!({
4409                "deviceId": "shared-peer",
4410                "ticket": "active-ticket",
4411                "topologyRole": "active",
4412                "topologyRevision": 3,
4413            })];
4414            registry
4415                .registrations
4416                .get_mut("space:backup")
4417                .unwrap()
4418                .desired_peers = vec![serde_json::json!({
4419                "deviceId": "shared-peer",
4420                "ticket": "backup-ticket",
4421                "topologyRole": "backup",
4422                "topologyRevision": 100,
4423            })];
4424            let (revision, peers_json) = registry.aggregate_desired_peers().unwrap();
4425            (revision, serde_json::from_str(&peers_json).unwrap())
4426        }
4427
4428        let forward = aggregate(false);
4429        let reverse = aggregate(true);
4430        assert_eq!(
4431            forward, reverse,
4432            "registration order is not lifecycle authority"
4433        );
4434        assert_eq!(forward.0, 1);
4435        assert_eq!(forward.1.len(), 1);
4436        assert_eq!(forward.1[0]["ticket"], "active-ticket");
4437        assert_eq!(forward.1[0]["topologyRole"], "active");
4438        assert_eq!(forward.1[0]["topologyRevision"], 1);
4439    }
4440
4441    #[test]
4442    fn native_capability_registry_keeps_physical_peer_until_all_references_withdraw() {
4443        let mut registry = NativeCapabilityRegistry::default();
4444        for (key, kind) in [("room:a", "room"), ("space:b", "space")] {
4445            registry
4446                .register(key.to_string(), kind.to_string(), key.to_string(), false)
4447                .unwrap();
4448        }
4449        registry
4450            .registrations
4451            .get_mut("room:a")
4452            .unwrap()
4453            .desired_peers = vec![serde_json::json!({
4454            "deviceId": "shared-peer",
4455            "topologyRole": "backup",
4456            "topologyRevision": 41,
4457        })];
4458        registry
4459            .registrations
4460            .get_mut("space:b")
4461            .unwrap()
4462            .desired_peers = vec![serde_json::json!({
4463            "deviceId": "shared-peer",
4464            "nodeId": "shared-node",
4465            "ticket": "space-ticket",
4466            "topologyRole": "active",
4467            "topologyRevision": 7,
4468        })];
4469
4470        let (_, both_json) = registry.aggregate_desired_peers().unwrap();
4471        let both: Vec<serde_json::Value> = serde_json::from_str(&both_json).unwrap();
4472        assert_eq!(both.len(), 1);
4473        assert_eq!(both[0]["topologyRole"], "active");
4474        assert_eq!(both[0]["ticket"], "space-ticket");
4475        assert_eq!(
4476            registry.registrations["room:a"].desired_peers[0]["topologyRevision"], 41,
4477            "root projection must not rewrite avenue-local lease state",
4478        );
4479
4480        registry
4481            .registrations
4482            .get_mut("space:b")
4483            .unwrap()
4484            .desired_peers
4485            .clear();
4486        let (_, room_only_json) = registry.aggregate_desired_peers().unwrap();
4487        let room_only: Vec<serde_json::Value> = serde_json::from_str(&room_only_json).unwrap();
4488        assert_eq!(
4489            room_only.len(),
4490            1,
4491            "one remaining capability keeps the leg desired"
4492        );
4493        assert_eq!(room_only[0]["topologyRole"], "backup");
4494
4495        registry
4496            .registrations
4497            .get_mut("room:a")
4498            .unwrap()
4499            .desired_peers
4500            .clear();
4501        let (_, empty_json) = registry.aggregate_desired_peers().unwrap();
4502        let empty: Vec<serde_json::Value> = serde_json::from_str(&empty_json).unwrap();
4503        assert!(
4504            empty.is_empty(),
4505            "physical desire ends only after every reference withdraws"
4506        );
4507    }
4508
4509    #[test]
4510    fn native_capability_registry_scopes_peer_data_and_rejects_unscoped_streams() {
4511        let mut registry = NativeCapabilityRegistry::default();
4512        registry
4513            .register(
4514                "devices:user-1".to_string(),
4515                "devices".to_string(),
4516                "user-1".to_string(),
4517                false,
4518            )
4519            .unwrap();
4520        registry
4521            .registrations
4522            .get_mut("devices:user-1")
4523            .unwrap()
4524            .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
4525
4526        let explicit = openrtc::client::NativePeerDataEvent {
4527            connection_id: "peer-1".to_string(),
4528            remote_node_id: None,
4529            transport: "webrtc".to_string(),
4530            transport_stable_id: 1,
4531            transport_generation: 1,
4532            route_generation: 1,
4533            payload: serde_json::to_vec(&serde_json::json!({
4534                "capability": "devices:user-1",
4535                "kind": "raw",
4536                "payload": [1, 2, 3]
4537            }))
4538            .unwrap(),
4539        };
4540        assert_eq!(
4541            registry.capability_keys_for_peer_data(&explicit),
4542            vec!["devices:user-1".to_string()]
4543        );
4544
4545        let unknown = openrtc::client::NativePeerDataEvent {
4546            payload: serde_json::to_vec(&serde_json::json!({
4547                "capability": "space:unknown",
4548                "kind": "raw"
4549            }))
4550            .unwrap(),
4551            ..explicit.clone()
4552        };
4553        assert!(registry.capability_keys_for_peer_data(&unknown).is_empty());
4554
4555        registry
4556            .register(
4557                "space:shared".to_string(),
4558                "space".to_string(),
4559                "shared".to_string(),
4560                false,
4561            )
4562            .unwrap();
4563        registry
4564            .registrations
4565            .get_mut("space:shared")
4566            .unwrap()
4567            .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
4568        let unscoped = openrtc::client::NativePeerDataEvent {
4569            payload: serde_json::to_vec(&serde_json::json!({
4570                "kind": "raw",
4571                "payload": [1, 2, 3]
4572            }))
4573            .unwrap(),
4574            ..explicit.clone()
4575        };
4576        assert!(
4577            registry.capability_keys_for_peer_data(&unscoped).is_empty(),
4578            "a shared physical peer cannot fan unscoped data into two avenues",
4579        );
4580
4581        let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
4582            channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
4583            metadata: None,
4584        };
4585        assert_eq!(
4586            registry.projected_stream_capability(&unscoped_channel),
4587            None,
4588            "stream ownership must never be inferred from capability count"
4589        );
4590    }
4591
4592    #[test]
4593    fn native_device_proof_uses_host_signer_without_exporting_private_key() {
4594        let state = OpenRtcTauriState::new(test_config())
4595            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4596        assert_eq!(
4597            state
4598                .native_device_key_signer
4599                .as_deref()
4600                .unwrap()
4601                .offline_assurance("app_native_test"),
4602            openrtc::offline::OfflineAssurance::Software
4603        );
4604        let public = state
4605            .device_public_key("app_native_test")
4606            .expect("public key");
4607        assert_eq!(public["kty"], "OKP");
4608        assert_eq!(public["crv"], "Ed25519");
4609        assert!(public.get("d").is_none());
4610        let signature = state
4611            .sign_device_proof("app_native_test", "openrtc:v2:test")
4612            .expect("signature");
4613        assert_eq!(
4614            base64::engine::general_purpose::URL_SAFE_NO_PAD
4615                .decode(signature)
4616                .expect("base64"),
4617            vec![7; 64]
4618        );
4619    }
4620
4621    #[test]
4622    fn offline_support_reports_signer_and_compiled_lan_truth() {
4623        let unavailable = OpenRtcTauriState::new(test_config()).offline_runtime_support();
4624        assert!(!unavailable.provisioning);
4625        assert_eq!(unavailable.local_mesh, cfg!(feature = "transport-lan"));
4626        assert_eq!(
4627            unavailable.reason,
4628            Some("native host device signer is unavailable")
4629        );
4630
4631        let available = OpenRtcTauriState::new(test_config())
4632            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()))
4633            .offline_runtime_support();
4634        assert!(available.provisioning);
4635        assert_eq!(available.local_mesh, cfg!(feature = "transport-lan"));
4636        assert_eq!(
4637            available.reason,
4638            (!cfg!(feature = "transport-lan"))
4639                .then_some("native host was built without transport-lan"),
4640        );
4641    }
4642
4643    #[tokio::test]
4644    async fn public_runtime_status_uses_host_facts_and_keeps_broadcast_unavailable() {
4645        use openrtc::client::CapabilityMaturity;
4646
4647        let unavailable_state = OpenRtcTauriState::new(test_config());
4648        let unavailable = unavailable_state
4649            .project_public_runtime_status(unavailable_state.client().runtime_status().await);
4650        assert_eq!(
4651            unavailable.product_maturity.offline_edge,
4652            CapabilityMaturity::Unavailable
4653        );
4654        assert_eq!(
4655            unavailable.product_maturity.broadcast,
4656            CapabilityMaturity::Unavailable
4657        );
4658        let serialized = serde_json::to_value(&unavailable).expect("serialized runtime status");
4659        assert_eq!(serialized["productMaturity"]["offlineEdge"], "unavailable");
4660        assert_eq!(serialized["productMaturity"]["broadcast"], "unavailable");
4661
4662        let installed_state = OpenRtcTauriState::new(test_config())
4663            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4664        let installed = installed_state
4665            .project_public_runtime_status(installed_state.client().runtime_status().await);
4666        assert_eq!(
4667            installed.product_maturity.offline_edge,
4668            if cfg!(feature = "transport-lan") {
4669                CapabilityMaturity::Preview
4670            } else {
4671                CapabilityMaturity::SupportOnly
4672            }
4673        );
4674        assert_eq!(
4675            installed.product_maturity.broadcast,
4676            CapabilityMaturity::Unavailable
4677        );
4678    }
4679
4680    #[test]
4681    fn native_certificate_records_round_trip_through_host_secure_storage() {
4682        let state = OpenRtcTauriState::new(test_config())
4683            .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
4684        let app_tag = "app_native_test";
4685        let key = "openrtc:v2:device-session:abc:device-1";
4686        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
4687        state
4688            .write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
4689            .unwrap();
4690        assert_eq!(
4691            state.read_secure_record(app_tag, key).unwrap().as_deref(),
4692            Some("{\"token\":\"bound\"}")
4693        );
4694        state.delete_secure_record(app_tag, key).unwrap();
4695        assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
4696    }
4697
4698    #[test]
4699    fn native_device_proof_fails_closed_without_secure_host_signer() {
4700        let state = OpenRtcTauriState::new(test_config());
4701        assert!(state
4702            .device_public_key("app_native_test")
4703            .expect_err("missing signer must fail")
4704            .contains("secure-storage signer"));
4705    }
4706
4707    fn native_transport_install_context() -> InstallContext {
4708        InstallContext {
4709            data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
4710        }
4711    }
4712
4713    fn managed_session_test_result(device_id: &str) -> StartSessionResult {
4714        StartSessionResult {
4715            local_node_id: format!("node-{device_id}"),
4716            ticket_scope: Some("user-device".to_string()),
4717            ticket: None,
4718            presence_started: true,
4719            auto_connect_started: true,
4720            local_device: openrtc::native_device::NativeDeviceIdentity {
4721                device_id: device_id.to_string(),
4722                device_name: "Test Device".to_string(),
4723                created_at_ms: 1,
4724                updated_at_ms: 1,
4725                name_source: None,
4726                system_info: None,
4727            },
4728        }
4729    }
4730
4731    async fn simulate_managed_session_start(
4732        state: Arc<OpenRtcTauriState>,
4733        device_id: String,
4734        pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
4735        starts: Arc<AtomicUsize>,
4736        stopped_owners: Arc<Mutex<Vec<usize>>>,
4737    ) -> SessionDisposition {
4738        let _start_guard = state.managed_session_start_guard.lock().await;
4739        let client = state.client();
4740        let key = format!("{}:{device_id}", client.app_tag());
4741        let (disposition, active_epoch) = {
4742            let active = state.managed_session.lock().await;
4743            let disposition = managed_session_disposition(
4744                active.as_ref().map(|record| record.key.as_str()),
4745                active
4746                    .as_ref()
4747                    .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
4748                active
4749                    .as_ref()
4750                    .is_some_and(|record| record.result.presence_started),
4751                active
4752                    .as_ref()
4753                    .is_some_and(|record| record.result.auto_connect_started),
4754                &key,
4755                true,
4756                true,
4757            );
4758            (
4759                disposition,
4760                active.as_ref().map(|record| record.owner_epoch),
4761            )
4762        };
4763
4764        if disposition == SessionDisposition::Reuse {
4765            return disposition;
4766        }
4767
4768        if disposition == SessionDisposition::Replace {
4769            let previous = state
4770                .managed_session
4771                .lock()
4772                .await
4773                .take()
4774                .expect("replacement must have a previous owner");
4775            stopped_owners
4776                .lock()
4777                .await
4778                .push(Arc::as_ptr(&previous.owner_client) as usize);
4779        }
4780
4781        starts.fetch_add(1, Ordering::SeqCst);
4782        if let Some((entered, release)) = pause {
4783            entered.send(()).expect("start observer must be waiting");
4784            release.await.expect("start release must be sent");
4785        }
4786
4787        let owner_epoch = match disposition {
4788            SessionDisposition::Refresh => {
4789                active_epoch.expect("refresh must preserve the active owner epoch")
4790            }
4791            SessionDisposition::Start | SessionDisposition::Replace => {
4792                state.allocate_managed_session_owner_epoch()
4793            }
4794            SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
4795        };
4796        let result = managed_session_test_result(&device_id);
4797        *state.managed_session.lock().await = Some(ManagedSessionRecord {
4798            key,
4799            owner_client: client,
4800            owner_epoch,
4801            result,
4802        });
4803        disposition
4804    }
4805
4806    #[tokio::test]
4807    async fn concurrent_same_key_start_reuses_one_owner_epoch() {
4808        let state = Arc::new(OpenRtcTauriState::new(test_config()));
4809        let starts = Arc::new(AtomicUsize::new(0));
4810        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
4811        let (first_entered_tx, first_entered_rx) = oneshot::channel();
4812        let (first_release_tx, first_release_rx) = oneshot::channel();
4813
4814        let first = tokio::spawn(simulate_managed_session_start(
4815            state.clone(),
4816            "same-device".to_string(),
4817            Some((first_entered_tx, first_release_rx)),
4818            starts.clone(),
4819            stopped_owners.clone(),
4820        ));
4821        first_entered_rx.await.expect("first start must pause");
4822
4823        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
4824        let second_state = state.clone();
4825        let second_starts = starts.clone();
4826        let second_stopped_owners = stopped_owners.clone();
4827        let second = tokio::spawn(async move {
4828            second_attempting_tx.send(()).unwrap();
4829            simulate_managed_session_start(
4830                second_state,
4831                "same-device".to_string(),
4832                None,
4833                second_starts,
4834                second_stopped_owners,
4835            )
4836            .await
4837        });
4838        second_attempting_rx.await.unwrap();
4839        tokio::task::yield_now().await;
4840        assert_eq!(starts.load(Ordering::SeqCst), 1);
4841
4842        first_release_tx.send(()).unwrap();
4843        assert_eq!(
4844            first.await.expect("first start task"),
4845            SessionDisposition::Start
4846        );
4847        assert_eq!(
4848            second.await.expect("second start task"),
4849            SessionDisposition::Reuse
4850        );
4851        assert_eq!(starts.load(Ordering::SeqCst), 1);
4852
4853        let active = state
4854            .managed_session
4855            .lock()
4856            .await
4857            .clone()
4858            .expect("active owner");
4859        assert_eq!(active.owner_epoch, 1);
4860        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
4861        assert!(stopped_owners.lock().await.is_empty());
4862    }
4863
4864    #[tokio::test]
4865    async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
4866        let state = Arc::new(OpenRtcTauriState::new(test_config()));
4867        let starts = Arc::new(AtomicUsize::new(0));
4868        let stopped_owners = Arc::new(Mutex::new(Vec::new()));
4869        let (first_entered_tx, first_entered_rx) = oneshot::channel();
4870        let (first_release_tx, first_release_rx) = oneshot::channel();
4871
4872        let first = tokio::spawn(simulate_managed_session_start(
4873            state.clone(),
4874            "first-device".to_string(),
4875            Some((first_entered_tx, first_release_rx)),
4876            starts.clone(),
4877            stopped_owners.clone(),
4878        ));
4879        first_entered_rx.await.expect("first start must pause");
4880        let first_owner = state.client();
4881
4882        let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
4883        let second_state = state.clone();
4884        let second_starts = starts.clone();
4885        let second_stopped_owners = stopped_owners.clone();
4886        let second = tokio::spawn(async move {
4887            second_attempting_tx.send(()).unwrap();
4888            simulate_managed_session_start(
4889                second_state,
4890                "second-device".to_string(),
4891                None,
4892                second_starts,
4893                second_stopped_owners,
4894            )
4895            .await
4896        });
4897        second_attempting_rx.await.unwrap();
4898        tokio::task::yield_now().await;
4899        assert!(Arc::ptr_eq(&first_owner, &state.client()));
4900
4901        first_release_tx.send(()).unwrap();
4902        assert_eq!(
4903            first.await.expect("first start task"),
4904            SessionDisposition::Start
4905        );
4906        assert_eq!(
4907            second.await.expect("replacement start task"),
4908            SessionDisposition::Replace
4909        );
4910        assert_eq!(starts.load(Ordering::SeqCst), 2);
4911
4912        let active = state
4913            .managed_session
4914            .lock()
4915            .await
4916            .clone()
4917            .expect("active owner");
4918        assert_eq!(active.owner_epoch, 2);
4919        assert!(active.key.contains("second-device"));
4920        assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
4921        assert_eq!(
4922            stopped_owners.lock().await.as_slice(),
4923            [Arc::as_ptr(&first_owner) as usize]
4924        );
4925        assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
4926    }
4927
4928    #[test]
4929    fn derives_app_tag_from_api_key_by_default() {
4930        let config = OpenRtcTauriConfig {
4931            api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
4932            ..OpenRtcTauriConfig::default()
4933        };
4934
4935        assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
4936    }
4937
4938    #[tokio::test]
4939    async fn requested_native_transport_without_installer_keeps_base_route_available() {
4940        let state = OpenRtcTauriState::new(test_config());
4941        let client = state.client();
4942        ensure_requested_native_transports(
4943            &state,
4944            &client,
4945            Some(&requested_ble_config()),
4946            &native_transport_install_context(),
4947        )
4948        .await
4949        .expect("an unavailable optional transport must not fail base Iroh startup");
4950
4951        assert!(state.installed_native_transports.lock().await.is_empty());
4952    }
4953
4954    #[tokio::test]
4955    async fn native_transport_installer_is_idempotent_for_one_client() {
4956        let installs = Arc::new(AtomicUsize::new(0));
4957        let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
4958            Arc::new(FakeTransportInstaller {
4959                installs: installs.clone(),
4960            }),
4961        );
4962        let client = state.client();
4963        let config = requested_ble_config();
4964
4965        ensure_requested_native_transports(
4966            &state,
4967            &client,
4968            Some(&config),
4969            &native_transport_install_context(),
4970        )
4971        .await
4972        .expect("first install");
4973        ensure_requested_native_transports(
4974            &state,
4975            &client,
4976            Some(&config),
4977            &native_transport_install_context(),
4978        )
4979        .await
4980        .expect("idempotent install");
4981
4982        assert_eq!(installs.load(Ordering::SeqCst), 1);
4983    }
4984
4985    #[tokio::test]
4986    async fn unavailable_native_transport_keeps_base_route_available() {
4987        let state = OpenRtcTauriState::new(test_config())
4988            .with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
4989        let client = state.client();
4990
4991        ensure_requested_native_transports(
4992            &state,
4993            &client,
4994            Some(&requested_ble_config()),
4995            &native_transport_install_context(),
4996        )
4997        .await
4998        .expect("optional transport installation failure must not fail base Iroh startup");
4999
5000        assert!(state.installed_native_transports.lock().await.is_empty());
5001    }
5002
5003    #[test]
5004    fn invalid_or_missing_api_key_fails_closed() {
5005        assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
5006        let mut config = OpenRtcTauriConfig::default();
5007        config.api_key = "pk_test_short".to_string();
5008        assert!(config.validated_api_key().is_err());
5009    }
5010
5011    #[cfg(feature = "native-broadcast-moq")]
5012    #[test]
5013    fn draft16_broadcast_feature_installs_one_private_native_adapter() {
5014        let state = OpenRtcTauriState::new(test_config());
5015        assert!(state.native_broadcast_adapter.is_some());
5016        assert!(state.native_broadcast_moq_adapter.is_some());
5017    }
5018
5019    #[test]
5020    fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
5021        let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
5022        let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
5023        let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
5024        let api_key = format!("pk_test_{}", "c".repeat(40));
5025        std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
5026        std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
5027        std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
5028
5029        let config = OpenRtcTauriConfig::from_env();
5030
5031        assert_eq!(config.api_key, api_key);
5032        assert_eq!(
5033            config.app_tag().unwrap(),
5034            openrtc::app_tag_from_api_key(&api_key)
5035        );
5036
5037        match previous_api_key {
5038            Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
5039            None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
5040        }
5041        match previous_project {
5042            Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
5043            None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
5044        }
5045        match previous_app_tag {
5046            Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
5047            None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
5048        }
5049    }
5050
5051    #[test]
5052    fn native_state_has_one_constructor_owned_app_identity() {
5053        let config = test_config();
5054        let expected = openrtc::app_tag_from_api_key(&config.api_key);
5055        let state = OpenRtcTauriState::new(config);
5056        let first = state.client();
5057        let second = state.client();
5058
5059        assert_eq!(first.app_tag(), expected);
5060        assert!(Arc::ptr_eq(&first, &second));
5061    }
5062
5063    #[test]
5064    fn managed_session_uses_persisted_native_device_id() {
5065        let identity = openrtc::native_device::NativeDeviceIdentity {
5066            device_id: "persisted-native-device".to_string(),
5067            device_name: "Mac".to_string(),
5068            created_at_ms: 1,
5069            updated_at_ms: 1,
5070            name_source: None,
5071            system_info: None,
5072        };
5073
5074        assert_eq!(
5075            managed_session_device_id(Some(" desktop-e2e-native "), &identity),
5076            "persisted-native-device"
5077        );
5078        assert_eq!(
5079            managed_session_device_id(None, &identity),
5080            "persisted-native-device"
5081        );
5082    }
5083
5084    #[test]
5085    fn repeated_managed_session_triggers_converge_on_one_owner() {
5086        let key = "app:persisted-native-device";
5087        let cases = [
5088            ("react-remount", true, true, true, true),
5089            ("hmr", true, true, true, true),
5090            ("auth-refresh", true, true, true, true),
5091            ("resume", true, true, true, true),
5092            ("alias-change", true, true, true, true),
5093        ];
5094        for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
5095            assert_eq!(
5096                managed_session_disposition(
5097                    Some(key),
5098                    true,
5099                    active_presence,
5100                    active_auto,
5101                    key,
5102                    requested_presence,
5103                    requested_auto,
5104                ),
5105                SessionDisposition::Reuse,
5106                "{label} must reuse the authoritative tuple"
5107            );
5108        }
5109
5110        assert_eq!(
5111            managed_session_disposition(Some(key), true, false, true, key, true, true),
5112            SessionDisposition::Refresh,
5113            "failed presence startup must retry idempotently"
5114        );
5115        assert_eq!(
5116            managed_session_disposition(
5117                Some(key),
5118                true,
5119                true,
5120                true,
5121                "app:other-native-device",
5122                true,
5123                true,
5124            ),
5125            SessionDisposition::Replace,
5126            "a physical native device change must replace the previous lifecycle owner"
5127        );
5128        assert_eq!(
5129            managed_session_disposition(Some(key), false, true, true, key, true, true),
5130            SessionDisposition::Replace,
5131            "the same tuple on a different client epoch must replace the previous owner"
5132        );
5133    }
5134
5135    #[test]
5136    fn user_device_revocation_invalidates_the_cached_managed_ticket() {
5137        assert!(revokes_managed_session("user-device"));
5138        assert!(revokes_managed_session("  user-device  "));
5139        assert!(!revokes_managed_session("share:example"));
5140        assert!(!revokes_managed_session(""));
5141    }
5142
5143    #[test]
5144    fn managed_session_metadata_uses_authoritative_device_id() {
5145        let raw = serde_json::json!({
5146            "deviceId": "persisted-native-device",
5147            "assistantDevice": {
5148                "deviceId": "desktop-e2e-native"
5149            }
5150        })
5151        .to_string();
5152
5153        let metadata =
5154            metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
5155        let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
5156
5157        assert_eq!(parsed["deviceId"], "desktop-e2e-native");
5158        assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
5159    }
5160}