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