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