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