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