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