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