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