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 openrtc::application_crypto_streams::{PeerRecvStream, PeerSendStream};
8use serde::{Deserialize, Serialize};
9use tauri::ipc::{Channel, InvokeBody, Request, Response};
10use tauri::{Manager, Runtime};
11use tokio::sync::Mutex;
12
13const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
14const PEER_BI_STREAM_ID_HEADER: &str = "x-openrtc-stream-id";
15const PEER_ID_HEADER: &str = "x-openrtc-peer-id";
16const MEDIA_PUBLICATION_ID_HEADER: &str = "x-openrtc-publication-id";
17const MEDIA_TIMESTAMP_US_HEADER: &str = "x-openrtc-timestamp-us";
18const MEDIA_DURATION_US_HEADER: &str = "x-openrtc-duration-us";
19const MEDIA_KEYFRAME_HEADER: &str = "x-openrtc-keyframe";
20const MEDIA_DISCARDABLE_HEADER: &str = "x-openrtc-discardable";
21
22fn request_body_bytes(request: &Request<'_>) -> Result<Vec<u8>, String> {
23 match request.body() {
24 InvokeBody::Raw(bytes) => Ok(bytes.clone()),
25 InvokeBody::Json(json) => serde_json::from_value::<Vec<u8>>(json.clone())
28 .map_err(|error| format!("invalid binary IPC payload: {error}")),
29 }
30}
31
32fn required_request_header(request: &Request<'_>, name: &str) -> Result<String, String> {
33 request
34 .headers()
35 .get(name)
36 .and_then(|value| value.to_str().ok())
37 .map(str::trim)
38 .filter(|value| !value.is_empty())
39 .map(ToOwned::to_owned)
40 .ok_or_else(|| format!("{name} header is required"))
41}
42
43fn parsed_request_header<T>(request: &Request<'_>, name: &str) -> Result<T, String>
44where
45 T: std::str::FromStr,
46 T::Err: std::fmt::Display,
47{
48 required_request_header(request, name)?
49 .parse::<T>()
50 .map_err(|error| format!("invalid {name} header: {error}"))
51}
52
53pub type InstallFuture = std::pin::Pin<
54 Box<
55 dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
56 + Send,
57 >,
58>;
59
60#[derive(Debug, Clone)]
61pub struct InstallContext {
62 pub data_dir: PathBuf,
66}
67
68pub trait TransportInstaller: Send + Sync {
72 fn id(&self) -> &'static str;
73 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
74 fn install(
75 &self,
76 client: Arc<openrtc::client::Client>,
77 config: openrtc::client::TransportConfig,
78 context: InstallContext,
79 ) -> InstallFuture;
80}
81
82pub trait DeviceKeySigner: Send + Sync {
88 fn public_jwk(&self, app_tag: &str) -> Result<serde_json::Value, String>;
89 fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String>;
90 fn delete(&self, app_tag: &str) -> Result<(), String>;
91 fn read_secure_record(&self, _app_tag: &str, _key: &str) -> Result<Option<String>, String> {
92 Ok(None)
93 }
94 fn write_secure_record(&self, _app_tag: &str, _key: &str, _value: &str) -> Result<(), String> {
95 Err("OpenRTC 2.0 native certificate persistence requires a host secure store".to_string())
96 }
97 fn delete_secure_record(&self, _app_tag: &str, _key: &str) -> Result<(), String> {
98 Ok(())
99 }
100}
101
102struct InstalledNativeTransport {
103 client_ptr: usize,
104 _runtime: Box<dyn std::any::Any + Send + Sync>,
105}
106
107#[derive(Debug, Clone, Serialize, Deserialize)]
108#[serde(rename_all = "camelCase")]
109pub struct OpenRtcTauriConfig {
110 pub api_key: String,
114 #[serde(default)]
115 pub data_dir: Option<PathBuf>,
116 #[serde(default)]
117 pub transport_config: Option<openrtc::client::TransportConfig>,
118}
119
120impl Default for OpenRtcTauriConfig {
121 fn default() -> Self {
122 Self {
123 api_key: String::new(),
124 data_dir: None,
125 transport_config: None,
126 }
127 }
128}
129
130impl OpenRtcTauriConfig {
131 pub fn from_env() -> Self {
132 let api_key = first_env(&[
133 "VITE_OPENRTC_KEY",
134 "VITE_OPENRTC_API_KEY",
135 "VITE_PLUTO_OPENRTC_API_KEY",
136 "OPENRTC_API_KEY",
137 ]);
138 Self {
139 api_key: api_key.unwrap_or_default(),
140 data_dir: None,
141 transport_config: None,
142 }
143 }
144
145 pub fn validated_api_key(&self) -> Result<&str, String> {
146 openrtc::validate_api_key(&self.api_key)
147 .map_err(|error| format!("invalid OpenRTC 2.0 public API key: {error}"))
148 }
149
150 pub fn app_tag(&self) -> Result<String, String> {
151 self.validated_api_key().map(openrtc::app_tag_from_api_key)
152 }
153}
154
155fn first_env(names: &[&str]) -> Option<String> {
156 names.iter().find_map(|name| {
157 std::env::var(name)
158 .ok()
159 .map(|value| value.trim().to_string())
160 .filter(|value| !value.is_empty())
161 })
162}
163
164#[derive(Default)]
165struct TokenRelayState {
166 identity_credential: RwLock<Option<String>>,
167}
168
169impl TokenRelayState {
170 fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
171 let relay = self.clone();
172 Box::new(move || {
173 relay
174 .identity_credential
175 .read()
176 .ok()
177 .and_then(|guard| guard.clone())
178 })
179 }
180
181 fn set(&self, identity_credential: Option<String>) {
182 if let Ok(mut guard) = self.identity_credential.write() {
183 *guard = normalize_token(identity_credential);
184 }
185 }
186}
187
188fn normalize_token(value: Option<String>) -> Option<String> {
189 value
190 .map(|value| value.trim().to_string())
191 .filter(|value| !value.is_empty())
192}
193
194#[derive(Debug, Clone)]
195struct NativeCapabilityRegistration {
196 avenue_kind: String,
197 avenue_id: String,
198 desired_revision: u64,
199 desired_peers: Vec<serde_json::Value>,
200}
201
202#[derive(Debug, Default)]
203struct NativeCapabilityRegistry {
204 registrations: HashMap<String, NativeCapabilityRegistration>,
205 root_desired_revision: u64,
206}
207
208impl NativeCapabilityRegistry {
209 fn register(
210 &mut self,
211 capability_key: String,
212 avenue_kind: String,
213 avenue_id: String,
214 ) -> Result<(), String> {
215 if let Some(existing) = self.registrations.get(&capability_key) {
216 if existing.avenue_kind == avenue_kind && existing.avenue_id == avenue_id {
217 return Ok(());
218 }
219 return Err(format!(
220 "native capability key {capability_key} is already registered for another avenue"
221 ));
222 }
223 self.registrations.insert(
224 capability_key,
225 NativeCapabilityRegistration {
226 avenue_kind,
227 avenue_id,
228 desired_revision: 0,
229 desired_peers: Vec::new(),
230 },
231 );
232 Ok(())
233 }
234
235 fn capability_keys_for_identity(
236 &self,
237 connection_id: Option<&str>,
238 device_id: Option<&str>,
239 device_id_hint: Option<&str>,
240 remote_node_id: Option<&str>,
241 ) -> Vec<String> {
242 let identities = [connection_id, device_id, device_id_hint, remote_node_id]
243 .into_iter()
244 .flatten()
245 .map(str::trim)
246 .filter(|value| !value.is_empty())
247 .collect::<BTreeSet<_>>();
248 self.registrations
249 .iter()
250 .filter_map(|(key, registration)| {
251 registration
252 .desired_peers
253 .iter()
254 .any(|peer| {
255 ["connectionId", "deviceId", "nodeId"]
256 .into_iter()
257 .filter_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
258 .map(str::trim)
259 .any(|value| identities.contains(value))
260 })
261 .then(|| key.clone())
262 })
263 .collect::<BTreeSet<_>>()
264 .into_iter()
265 .collect()
266 }
267
268 fn capability_keys_for_state(&self, snapshot: &openrtc::client::StateSnapshot) -> Vec<String> {
269 self.capability_keys_for_identity(
270 Some(&snapshot.connection_id),
271 snapshot.device_id.as_deref(),
272 snapshot.device_id_hint.as_deref(),
273 snapshot.remote_node_id.as_deref(),
274 )
275 }
276
277 fn capability_keys_for_peer_data(
278 &self,
279 event: &openrtc::client::NativePeerDataEvent,
280 ) -> Vec<String> {
281 if let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) {
282 if let Some(capability) = value
283 .get("capability")
284 .and_then(serde_json::Value::as_str)
285 .map(str::trim)
286 .filter(|value| !value.is_empty())
287 {
288 if self.registrations.contains_key(capability) {
289 return vec![capability.to_string()];
290 }
291 return Vec::new();
292 }
293 }
294 self.capability_keys_for_identity(
295 Some(&event.connection_id),
296 None,
297 None,
298 event.remote_node_id.as_deref(),
299 )
300 }
301
302 fn projected_stream_capability(
303 &self,
304 channel: &openrtc::stream_metadata::ChannelMetadata,
305 ) -> Option<String> {
306 let explicit = channel
307 .metadata
308 .as_ref()
309 .and_then(|metadata| metadata.get("openrtcCapability"))
310 .and_then(serde_json::Value::as_str)
311 .map(str::trim)
312 .filter(|value| !value.is_empty());
313 match explicit {
314 Some(key) if self.registrations.contains_key(key) => Some(key.to_string()),
315 _ => None,
316 }
317 }
318
319 fn aggregate_desired_peers(&mut self) -> Result<(u64, String), String> {
320 let mut peers = std::collections::BTreeMap::<String, serde_json::Value>::new();
321 for registration in self.registrations.values() {
322 for peer in ®istration.desired_peers {
323 let identity = ["deviceId", "nodeId", "ticket"]
324 .into_iter()
325 .find_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
326 .map(str::trim)
327 .filter(|value| !value.is_empty())
328 .ok_or_else(|| "desired peer has no stable identity".to_string())?;
329 peers
330 .entry(identity.to_string())
331 .or_insert_with(|| peer.clone());
332 }
333 }
334 self.root_desired_revision = self.root_desired_revision.saturating_add(1);
335 let payload = serde_json::to_string(&peers.into_values().collect::<Vec<_>>())
336 .map_err(|error| format!("serialize aggregated desired peers: {error}"))?;
337 Ok((self.root_desired_revision, payload))
338 }
339}
340
341fn required_native_capability_part(value: String, label: &str) -> Result<String, String> {
342 let value = value.trim().to_string();
343 if value.is_empty()
344 || value.len() > 192
345 || !value
346 .chars()
347 .all(|character| character.is_ascii_alphanumeric() || "_.:@-".contains(character))
348 {
349 return Err(format!("native {label} is invalid"));
350 }
351 Ok(value)
352}
353
354pub struct OpenRtcTauriState {
355 client: RwLock<Arc<openrtc::client::Client>>,
356 config: OpenRtcTauriConfig,
357 token_relay: Arc<TokenRelayState>,
358 data_dir: Option<PathBuf>,
359 connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
360 peer_data_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
361 presence_loop_active: Mutex<bool>,
362 subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
363 native_projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
364 projected_stream_offered: AtomicU64,
365 projected_stream_decoded: AtomicU64,
366 projected_stream_projected: AtomicU64,
367 projected_stream_unhandled_non_channel: AtomicU64,
368 projected_stream_unhandled_other_channel: AtomicU64,
369 projected_stream_unhandled_no_subscription: AtomicU64,
370 projected_stream_failures: AtomicU64,
371 projected_stream_last_transport_stable_id: AtomicU64,
372 projected_stream_validation_checks: AtomicU64,
373 projected_stream_validation_rejections: AtomicU64,
374 projected_stream_last_validated_transport_stable_id: AtomicU64,
375 projected_stream_authorized: AtomicU64,
376 projected_stream_unauthorized: AtomicU64,
377 projected_stream_last_channel: Mutex<Option<String>>,
378 projected_stream_last_protocol: Mutex<Option<String>>,
379 projected_stream_last_connection_id: Mutex<Option<String>>,
380 projected_stream_last_remote_node_id: Mutex<Option<String>>,
381 peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
382 peer_uni_streams: Mutex<HashMap<String, Arc<Mutex<Option<PeerSendStream>>>>>,
383 portable_media: Mutex<openrtc::media::PortableMediaSession>,
384 managed_session_start_guard: Mutex<()>,
385 managed_session_next_owner_epoch: AtomicU64,
386 managed_session: Mutex<Option<ManagedSessionRecord>>,
387 capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
388 native_transport_installers: Vec<Arc<dyn TransportInstaller>>,
389 installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
390 native_device_key_signer: Option<Arc<dyn DeviceKeySigner>>,
391}
392
393impl OpenRtcTauriState {
394 pub fn new(config: OpenRtcTauriConfig) -> Self {
395 openrtc::ensure_rustls();
396
397 let token_relay = Arc::new(TokenRelayState::default());
398 let client = build_client(&config, &token_relay)
399 .expect("OpenRTC Tauri 2.0 requires a valid public API key");
400 let data_dir = config.data_dir.clone();
401
402 Self {
403 client: RwLock::new(client),
404 config,
405 token_relay,
406 data_dir,
407 connection_state_forwarder: std::sync::Mutex::new(None),
408 peer_data_forwarder: std::sync::Mutex::new(None),
409 presence_loop_active: Mutex::new(false),
410 subscriptions: Mutex::new(HashMap::new()),
411 native_projection_subscription: Arc::new(Mutex::new(None)),
412 projected_stream_offered: AtomicU64::new(0),
413 projected_stream_decoded: AtomicU64::new(0),
414 projected_stream_projected: AtomicU64::new(0),
415 projected_stream_unhandled_non_channel: AtomicU64::new(0),
416 projected_stream_unhandled_other_channel: AtomicU64::new(0),
417 projected_stream_unhandled_no_subscription: AtomicU64::new(0),
418 projected_stream_failures: AtomicU64::new(0),
419 projected_stream_last_transport_stable_id: AtomicU64::new(0),
420 projected_stream_validation_checks: AtomicU64::new(0),
421 projected_stream_validation_rejections: AtomicU64::new(0),
422 projected_stream_last_validated_transport_stable_id: AtomicU64::new(0),
423 projected_stream_authorized: AtomicU64::new(0),
424 projected_stream_unauthorized: AtomicU64::new(0),
425 projected_stream_last_channel: Mutex::new(None),
426 projected_stream_last_protocol: Mutex::new(None),
427 projected_stream_last_connection_id: Mutex::new(None),
428 projected_stream_last_remote_node_id: Mutex::new(None),
429 peer_bi_streams: Mutex::new(HashMap::new()),
430 peer_uni_streams: Mutex::new(HashMap::new()),
431 portable_media: Mutex::new(openrtc::media::PortableMediaSession::default()),
432 managed_session_start_guard: Mutex::new(()),
433 managed_session_next_owner_epoch: AtomicU64::new(0),
434 managed_session: Mutex::new(None),
435 capability_registry: Arc::new(Mutex::new(NativeCapabilityRegistry::default())),
436 native_transport_installers: Vec::new(),
437 installed_native_transports: Mutex::new(HashMap::new()),
438 native_device_key_signer: None,
439 }
440 }
441
442 pub fn with_native_transport_installer(
443 mut self,
444 installer: Arc<dyn TransportInstaller>,
445 ) -> Self {
446 self.native_transport_installers.push(installer);
447 self
448 }
449
450 pub fn with_native_device_key_signer(mut self, signer: Arc<dyn DeviceKeySigner>) -> Self {
451 self.native_device_key_signer = Some(signer);
452 self
453 }
454
455 fn device_public_key(&self, app_tag: &str) -> Result<serde_json::Value, String> {
456 let app_tag = required_app_tag(app_tag)?;
457 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
458 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
459 })?;
460 let value = signer.public_jwk(app_tag)?;
461 validate_public_device_jwk(&value)?;
462 Ok(value)
463 }
464
465 fn sign_device_proof(&self, app_tag: &str, challenge: &str) -> Result<String, String> {
466 let app_tag = required_app_tag(app_tag)?;
467 if challenge.is_empty() || challenge.len() > 2_048 {
468 return Err("OpenRTC device challenge is invalid".to_string());
469 }
470 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
471 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
472 })?;
473 let signature = signer.sign(app_tag, challenge.as_bytes())?;
474 if signature.len() != 64 {
475 return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
476 }
477 Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
478 }
479
480 fn delete_device_key(&self, app_tag: &str) -> Result<(), String> {
481 let app_tag = required_app_tag(app_tag)?;
482 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
483 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
484 })?;
485 signer.delete(app_tag)
486 }
487
488 fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
489 let app_tag = required_app_tag(app_tag)?;
490 let key = required_secure_record_key(key)?;
491 self.native_device_key_signer
492 .as_ref()
493 .ok_or_else(|| {
494 "OpenRTC 2.0 native certificate persistence requires a host secure store"
495 .to_string()
496 })?
497 .read_secure_record(app_tag, key)
498 }
499
500 fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
501 let app_tag = required_app_tag(app_tag)?;
502 let key = required_secure_record_key(key)?;
503 if value.is_empty() || value.len() > 16 * 1024 {
504 return Err("OpenRTC secure record is invalid".to_string());
505 }
506 self.native_device_key_signer
507 .as_ref()
508 .ok_or_else(|| {
509 "OpenRTC 2.0 native certificate persistence requires a host secure store"
510 .to_string()
511 })?
512 .write_secure_record(app_tag, key, value)
513 }
514
515 fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
516 let app_tag = required_app_tag(app_tag)?;
517 let key = required_secure_record_key(key)?;
518 self.native_device_key_signer
519 .as_ref()
520 .ok_or_else(|| {
521 "OpenRTC 2.0 native certificate persistence requires a host secure store"
522 .to_string()
523 })?
524 .delete_secure_record(app_tag, key)
525 }
526
527 pub fn client(&self) -> Arc<openrtc::client::Client> {
528 self.client
529 .read()
530 .map(|guard| guard.clone())
531 .expect("OpenRTC Tauri client state is poisoned")
532 }
533
534 pub fn set_identity_credential(&self, identity_credential: Option<String>) {
535 self.token_relay.set(identity_credential);
536 }
537
538 fn allocate_managed_session_owner_epoch(&self) -> u64 {
539 self.managed_session_next_owner_epoch
540 .fetch_add(1, Ordering::SeqCst)
541 + 1
542 }
543
544 fn replace_connection_state_forwarder(&self, client: Arc<openrtc::client::Client>) {
545 let capability_registry = self.capability_registry.clone();
546 let projection_subscription = self.native_projection_subscription.clone();
547 let next = tauri::async_runtime::spawn(async move {
548 forward_connection_state_events(client, capability_registry, projection_subscription)
549 .await;
550 });
551 if let Ok(mut guard) = self.connection_state_forwarder.lock() {
552 if let Some(previous) = guard.replace(next) {
553 previous.abort();
554 }
555 }
556 }
557
558 fn replace_peer_data_forwarder(&self, client: Arc<openrtc::client::Client>) {
559 let capability_registry = self.capability_registry.clone();
560 let projection_subscription = self.native_projection_subscription.clone();
561 let next = tauri::async_runtime::spawn(async move {
562 forward_peer_data_events(client, capability_registry, projection_subscription).await;
563 });
564 if let Ok(mut guard) = self.peer_data_forwarder.lock() {
565 if let Some(previous) = guard.replace(next) {
566 previous.abort();
567 }
568 }
569 }
570}
571
572fn build_client(
573 config: &OpenRtcTauriConfig,
574 token_relay: &Arc<TokenRelayState>,
575) -> Result<Arc<openrtc::client::Client>, String> {
576 let mut builder = openrtc::client::Client::builder(
577 config.validated_api_key()?.to_string(),
578 token_relay.token_provider(),
579 )
580 .map_err(|error| error.to_string())?;
581 if let Some(transport_config) = config.transport_config.clone() {
582 builder = builder.transport_config(transport_config);
583 }
584 Ok(Arc::new(builder.build()))
585}
586
587fn requested_local_device_id(value: Option<&str>) -> Option<String> {
588 value
589 .map(str::trim)
590 .filter(|value| !value.is_empty())
591 .map(ToOwned::to_owned)
592}
593
594fn managed_session_device_id(
595 _requested: Option<&str>,
596 native_identity: &openrtc::native_device::NativeDeviceIdentity,
597) -> String {
598 native_identity.device_id.clone()
602}
603
604#[cfg(test)]
605fn metadata_with_authoritative_device_id(
606 metadata: Option<String>,
607 device_id: &str,
608) -> Option<String> {
609 let device_id = device_id.trim();
610 if device_id.is_empty() {
611 return metadata;
612 }
613
614 let Some(raw_metadata) = metadata else {
615 return Some(serde_json::json!({ "deviceId": device_id }).to_string());
616 };
617
618 match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
619 Ok(serde_json::Value::Object(mut map)) => {
620 map.insert(
621 "deviceId".to_string(),
622 serde_json::Value::String(device_id.to_string()),
623 );
624 Some(serde_json::Value::Object(map).to_string())
625 }
626 _ => Some(
627 serde_json::json!({
628 "deviceId": device_id,
629 "metadata": raw_metadata,
630 })
631 .to_string(),
632 ),
633 }
634}
635
636struct PeerBiStreamHandle {
637 send: Arc<Mutex<Option<PeerSendStream>>>,
638 recv: Option<PeerRecvStream>,
639 read_task: Option<tokio::task::JoinHandle<()>>,
640}
641
642struct NativeProjectionSubscription {
643 request_id: String,
644 channel: Channel<NativeProjectionEvent>,
645}
646
647#[derive(Serialize, Clone)]
648#[serde(rename_all = "camelCase")]
649struct NativeProjectionEvent {
650 request_id: String,
651 kind: &'static str,
652 capability_keys: Vec<String>,
653 payload: serde_json::Value,
654}
655
656#[derive(Serialize)]
657#[serde(rename_all = "camelCase")]
658pub struct OpenBiResult {
659 stream_id: String,
660 connection_id: Option<String>,
661 remote_node_id: String,
662}
663
664#[derive(Serialize)]
665#[serde(rename_all = "camelCase")]
666pub struct OpenUniResult {
667 stream_id: String,
668 connection_id: Option<String>,
669 remote_node_id: String,
670}
671
672#[derive(Serialize, Clone)]
673#[serde(rename_all = "camelCase")]
674struct IncomingPeerBiStreamEvent {
675 request_id: String,
676 stream_id: String,
677 connection_id: Option<String>,
678 remote_node_id: String,
679 transport_stable_id: u64,
680 channel: Option<openrtc::stream_metadata::ChannelMetadata>,
681 application_authorized: bool,
682 capability_key: String,
683}
684
685#[derive(Debug, Clone, Serialize)]
686#[serde(rename_all = "camelCase")]
687struct ProjectedStreamDiagnostics {
688 subscription_active: bool,
689 offered: u64,
690 decoded: u64,
691 projected: u64,
692 unhandled_non_channel: u64,
693 unhandled_other_channel: u64,
694 unhandled_no_subscription: u64,
695 failures: u64,
696 last_transport_stable_id: u64,
697 validation_checks: u64,
698 validation_rejections: u64,
699 last_validated_transport_stable_id: u64,
700 authorized: u64,
701 unauthorized: u64,
702 last_channel: Option<String>,
703 last_protocol: Option<String>,
704 last_connection_id: Option<String>,
705 last_remote_node_id: Option<String>,
706}
707
708pub enum ProjectedPeerBiHandoff {
713 Handled,
714 Unhandled {
715 send: PeerSendStream,
716 recv: PeerRecvStream,
717 },
718}
719
720#[derive(Serialize, Clone)]
721#[serde(rename_all = "camelCase")]
722struct PeerBiStreamClosedEvent {
723 r#type: &'static str,
724 error: Option<String>,
725}
726
727#[derive(Debug, Clone, Serialize)]
728#[serde(rename_all = "camelCase")]
729pub struct StartSessionResult {
730 local_node_id: String,
731 ticket_scope: Option<String>,
732 ticket: Option<String>,
733 presence_started: bool,
734 auto_connect_started: bool,
735 local_device: openrtc::native_device::NativeDeviceIdentity,
736}
737
738#[derive(Clone)]
739struct ManagedSessionRecord {
740 key: String,
741 owner_client: Arc<openrtc::client::Client>,
742 owner_epoch: u64,
743 result: StartSessionResult,
744}
745
746#[derive(Debug, Clone, Copy, PartialEq, Eq)]
747enum SessionDisposition {
748 Start,
749 Reuse,
750 Refresh,
751 Replace,
752}
753
754fn managed_session_disposition(
755 active_key: Option<&str>,
756 active_client_matches: bool,
757 active_presence_started: bool,
758 active_auto_connect_started: bool,
759 requested_key: &str,
760 requested_presence: bool,
761 requested_auto_connect: bool,
762) -> SessionDisposition {
763 let Some(active_key) = active_key else {
764 return SessionDisposition::Start;
765 };
766 if !active_client_matches || active_key != requested_key {
767 return SessionDisposition::Replace;
768 }
769 if (!requested_presence || active_presence_started)
770 && (!requested_auto_connect || active_auto_connect_started)
771 {
772 SessionDisposition::Reuse
773 } else {
774 SessionDisposition::Refresh
775 }
776}
777
778fn revokes_managed_session(scope: &str) -> bool {
779 scope.trim() == "user-device"
780}
781
782pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
783 init_with_state(OpenRtcTauriState::new(config))
784}
785
786pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
787 tauri::plugin::Builder::new(PLUGIN_NAME)
788 .setup(move |app, _api| {
789 let client = state.client();
790 let transport_config = state.config.transport_config.clone();
791 let install_context = InstallContext {
796 data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
797 };
798 tauri::async_runtime::block_on(ensure_requested_native_transports(
799 &state,
800 &client,
801 transport_config.as_ref(),
802 &install_context,
803 ))
804 .map_err(std::io::Error::other)?;
805 state.replace_connection_state_forwarder(client.clone());
806 state.replace_peer_data_forwarder(client);
807 app.manage(state);
808 Ok(())
809 })
810 .invoke_handler(tauri::generate_handler![
811 openrtc_set_identity_credential,
812 openrtc_device_public_key,
813 openrtc_sign_device_proof,
814 openrtc_delete_device_key,
815 openrtc_read_secure_record,
816 openrtc_write_secure_record,
817 openrtc_delete_secure_record,
818 rtc_native_status,
819 get_rtc_local_device_info,
820 update_rtc_local_device_name,
821 get_iroh_node_id,
822 start_iroh_node,
823 get_iroh_endpoint_ticket,
824 register_session_token,
825 get_endpoint_ticket_with_token,
826 validate_session_token,
827 revoke_session_tokens_by_scope,
828 stop_rtc_presence_loop,
829 register_rtc_capability,
830 unregister_rtc_capability,
831 start_rtc_managed_session,
832 start_rtc_external_auto_connect,
833 submit_rtc_desired_peers,
834 stop_rtc_auto_connect,
835 notify_rtc_network_change,
836 set_rtc_transport_priority,
837 connect_to_device,
838 disconnect_device,
839 set_auto_connect_excluded,
840 set_rtc_external_auto_connect_excluded,
841 resolve_rtc_peer_connection_records,
842 resolve_rtc_peer_identity,
843 get_rtc_peer_session,
844 list_rtc_peer_sessions,
845 list_rtc_managed_connections,
846 wait_for_rtc_settled_peer,
847 list_rtc_connection_states,
848 get_rtc_connection_state,
849 stop_rtc_subscription,
850 start_native_projection,
851 get_projected_stream_diagnostics,
852 is_current_transport_stable_id,
853 send_peer_message,
854 is_peer_connected,
855 open_peer_bi_stream,
856 open_peer_bi_transport_only_stream,
857 open_peer_uni_stream,
858 write_peer_bi_stream,
859 start_peer_bi_stream_read,
860 cancel_peer_bi_stream_read,
861 close_peer_bi_stream,
862 write_peer_uni_stream,
863 close_peer_uni_stream,
864 begin_openrtc_media_publication,
865 encode_openrtc_media_sample,
866 pause_openrtc_media_publication,
867 retire_openrtc_media_publication,
868 retire_openrtc_media_receiver,
869 decode_openrtc_media_chunk,
870 decode_openrtc_media_control,
871 ])
872 .build()
873}
874
875pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
876 init(OpenRtcTauriConfig::from_env())
877}
878
879async fn send_native_projection<T: Serialize>(
880 subscription: &Arc<Mutex<Option<NativeProjectionSubscription>>>,
881 kind: &'static str,
882 capability_keys: Vec<String>,
883 payload: &T,
884) {
885 if capability_keys.is_empty() {
886 return;
887 }
888 let Ok(payload) = serde_json::to_value(payload) else {
889 return;
890 };
891 let target = subscription.lock().await.as_ref().map(|subscription| {
892 (
893 subscription.request_id.clone(),
894 subscription.channel.clone(),
895 )
896 });
897 let Some((request_id, channel)) = target else {
898 return;
899 };
900 let _ = channel.send(NativeProjectionEvent {
901 request_id,
902 kind,
903 capability_keys,
904 payload,
905 });
906}
907
908async fn forward_connection_state_events(
909 client: Arc<openrtc::client::Client>,
910 capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
911 projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
912) {
913 let mut rx = client.connection_state_updates();
914 while let Ok(snapshot) = rx.recv().await {
915 let keys = capability_registry
916 .lock()
917 .await
918 .capability_keys_for_state(&snapshot);
919 send_native_projection(
920 &projection_subscription,
921 "connection-state",
922 keys,
923 &snapshot,
924 )
925 .await;
926 }
927}
928
929async fn forward_peer_data_events(
930 client: Arc<openrtc::client::Client>,
931 capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
932 projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
933) {
934 let mut rx = client.subscribe_native_peer_data();
935 while let Ok(event) = rx.recv().await {
936 let keys = capability_registry
937 .lock()
938 .await
939 .capability_keys_for_peer_data(&event);
940 send_native_projection(&projection_subscription, "peer-data", keys, &event).await;
941 }
942}
943
944async fn replay_current_connection_states(state: &OpenRtcTauriState) {
945 let snapshots = state.client().connection_states().await;
946 let registry = state.capability_registry.lock().await;
947 for snapshot in snapshots {
948 let keys = registry.capability_keys_for_state(&snapshot);
949 send_native_projection(
950 &state.native_projection_subscription,
951 "connection-state",
952 keys,
953 &snapshot,
954 )
955 .await;
956 }
957}
958
959fn app_data_dir<R: Runtime>(
960 app: &tauri::AppHandle<R>,
961 state: &OpenRtcTauriState,
962) -> Result<PathBuf, String> {
963 state
964 .data_dir
965 .clone()
966 .or_else(|| app.path().app_data_dir().ok())
967 .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
968}
969
970async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
971 if let Some(node_id) = client.current_node_id().await {
972 return Ok(node_id);
973 }
974 client
975 .init_iroh(None, Vec::new())
976 .await
977 .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
978}
979
980async fn ensure_requested_native_transports(
981 state: &OpenRtcTauriState,
982 client: &Arc<openrtc::client::Client>,
983 transports: Option<&openrtc::client::TransportConfig>,
984 context: &InstallContext,
985) -> Result<(), String> {
986 let Some(transports) = transports else {
987 return Ok(());
988 };
989 let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
990 if ble_requested
991 && !state
992 .native_transport_installers
993 .iter()
994 .any(|installer| installer.is_requested(transports))
995 {
996 eprintln!(
997 "[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
998 );
999 return Ok(());
1000 }
1001 let client_ptr = Arc::as_ptr(client) as usize;
1002 for installer in &state.native_transport_installers {
1003 if !installer.is_requested(transports) {
1004 continue;
1005 }
1006 let already_installed = state
1007 .installed_native_transports
1008 .lock()
1009 .await
1010 .get(installer.id())
1011 .is_some_and(|installed| installed.client_ptr == client_ptr);
1012 if already_installed {
1013 continue;
1014 }
1015 if client.current_node_id().await.is_some() {
1016 eprintln!(
1017 "[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
1018 installer.id()
1019 );
1020 continue;
1021 }
1022 let runtime = match installer
1023 .install(client.clone(), transports.clone(), context.clone())
1024 .await
1025 {
1026 Ok(runtime) => runtime,
1027 Err(error) => {
1028 eprintln!(
1029 "[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
1030 installer.id(),
1031 error
1032 );
1033 continue;
1034 }
1035 };
1036 state.installed_native_transports.lock().await.insert(
1037 installer.id(),
1038 InstalledNativeTransport {
1039 client_ptr,
1040 _runtime: runtime,
1041 },
1042 );
1043 }
1044 Ok(())
1045}
1046
1047async fn register_peer_bi_stream(
1048 state: &OpenRtcTauriState,
1049 connection_id: Option<String>,
1050 remote_node_id: String,
1051 send: PeerSendStream,
1052 recv: PeerRecvStream,
1053) -> OpenBiResult {
1054 let stream_id = uuid::Uuid::new_v4().to_string();
1055 let send = Arc::new(Mutex::new(Some(send)));
1056
1057 state.peer_bi_streams.lock().await.insert(
1058 stream_id.clone(),
1059 PeerBiStreamHandle {
1060 send,
1061 recv: Some(recv),
1064 read_task: None,
1065 },
1066 );
1067
1068 OpenBiResult {
1069 stream_id,
1070 connection_id,
1071 remote_node_id,
1072 }
1073}
1074
1075async fn read_peer_exact(recv: &mut PeerRecvStream, buffer: &mut [u8]) -> std::io::Result<()> {
1076 let mut offset = 0;
1077 while offset < buffer.len() {
1078 let read = recv.read(&mut buffer[offset..]).await?;
1079 if read == 0 {
1080 return Err(std::io::Error::new(
1081 std::io::ErrorKind::UnexpectedEof,
1082 "peer stream ended during channel classification",
1083 ));
1084 }
1085 offset += read;
1086 }
1087 Ok(())
1088}
1089
1090fn is_openrtc_projected_channel(channel: &openrtc::stream_metadata::ChannelMetadata) -> bool {
1091 channel.channel_id == openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID
1092}
1093
1094pub async fn handoff_projected_peer_bi_stream<R: Runtime>(
1102 app: tauri::AppHandle<R>,
1103 connection_id: Option<String>,
1104 remote_node_id: String,
1105 transport_stable_id: u64,
1106 send: PeerSendStream,
1107 mut recv: PeerRecvStream,
1108) -> Result<ProjectedPeerBiHandoff, String> {
1109 let state = app.state::<OpenRtcTauriState>();
1110 state
1111 .projected_stream_offered
1112 .fetch_add(1, Ordering::Relaxed);
1113 state
1114 .projected_stream_last_transport_stable_id
1115 .store(transport_stable_id, Ordering::Relaxed);
1116 *state.projected_stream_last_connection_id.lock().await = connection_id.clone();
1117 *state.projected_stream_last_remote_node_id.lock().await = Some(remote_node_id.clone());
1118 let mut envelope = Vec::new();
1119 let channel = loop {
1120 match openrtc::stream_metadata::decode_prefix(&envelope) {
1121 openrtc::stream_metadata::DecodeDecision::NotMatched => {
1122 state
1123 .projected_stream_unhandled_non_channel
1124 .fetch_add(1, Ordering::Relaxed);
1125 return Ok(ProjectedPeerBiHandoff::Unhandled {
1126 send,
1127 recv: recv.with_plaintext_prefix(envelope),
1128 });
1129 }
1130 openrtc::stream_metadata::DecodeDecision::Decoded(prefix) => break prefix.channel,
1131 openrtc::stream_metadata::DecodeDecision::NeedMore(required) => {
1132 let missing = required.saturating_sub(envelope.len());
1133 if missing == 0 {
1134 return Err("channel-envelope decoder made no progress".to_string());
1135 }
1136 let mut bytes = vec![0u8; missing];
1137 if let Err(error) = read_peer_exact(&mut recv, &mut bytes).await {
1138 state
1139 .projected_stream_failures
1140 .fetch_add(1, Ordering::Relaxed);
1141 return Err(format!("read projected channel envelope: {error}"));
1142 }
1143 envelope.extend_from_slice(&bytes);
1144 }
1145 }
1146 };
1147 state
1148 .projected_stream_decoded
1149 .fetch_add(1, Ordering::Relaxed);
1150 *state.projected_stream_last_channel.lock().await = Some(channel.channel_id.clone());
1151 *state.projected_stream_last_protocol.lock().await = channel
1152 .metadata
1153 .as_ref()
1154 .and_then(|metadata| metadata.get("protocol"))
1155 .and_then(serde_json::Value::as_str)
1156 .map(str::to_string);
1157
1158 if !is_openrtc_projected_channel(&channel) {
1159 state
1160 .projected_stream_unhandled_other_channel
1161 .fetch_add(1, Ordering::Relaxed);
1162 return Ok(ProjectedPeerBiHandoff::Unhandled {
1163 send,
1164 recv: recv.with_plaintext_prefix(envelope),
1165 });
1166 }
1167
1168 let Some(capability_key) = state
1169 .capability_registry
1170 .lock()
1171 .await
1172 .projected_stream_capability(&channel)
1173 else {
1174 state
1178 .projected_stream_unhandled_other_channel
1179 .fetch_add(1, Ordering::Relaxed);
1180 return Ok(ProjectedPeerBiHandoff::Unhandled {
1181 send,
1182 recv: recv.with_plaintext_prefix(envelope),
1183 });
1184 };
1185
1186 let subscription = state.native_projection_subscription.lock().await;
1187 let Some((request_id, projection_channel)) = subscription.as_ref().map(|subscription| {
1188 (
1189 subscription.request_id.clone(),
1190 subscription.channel.clone(),
1191 )
1192 }) else {
1193 state
1194 .projected_stream_unhandled_no_subscription
1195 .fetch_add(1, Ordering::Relaxed);
1196 return Ok(ProjectedPeerBiHandoff::Unhandled {
1197 send,
1198 recv: recv.with_plaintext_prefix(envelope),
1199 });
1200 };
1201 drop(subscription);
1202
1203 let application_authorized = connection_id.as_deref().is_some_and(|connection_id| {
1204 matches!(
1205 state.client().session_admission(connection_id),
1206 openrtc::session_token::SessionAdmission::Accepted { .. }
1207 )
1208 });
1209 if application_authorized {
1210 state
1211 .projected_stream_authorized
1212 .fetch_add(1, Ordering::Relaxed);
1213 } else {
1214 state
1215 .projected_stream_unauthorized
1216 .fetch_add(1, Ordering::Relaxed);
1217 }
1218
1219 let result = register_peer_bi_stream(
1220 state.inner(),
1221 connection_id,
1222 remote_node_id.clone(),
1223 send,
1224 recv,
1225 )
1226 .await;
1227 let event = IncomingPeerBiStreamEvent {
1228 request_id: request_id.clone(),
1229 stream_id: result.stream_id.clone(),
1230 connection_id: result.connection_id,
1231 remote_node_id,
1232 transport_stable_id,
1233 channel: Some(channel),
1234 application_authorized,
1235 capability_key: capability_key.clone(),
1236 };
1237 let payload = serde_json::to_value(event)
1238 .map_err(|error| format!("serialize projected peer stream: {error}"))?;
1239 if let Err(error) = projection_channel.send(NativeProjectionEvent {
1240 request_id,
1241 kind: "stream",
1242 capability_keys: vec![capability_key],
1243 payload,
1244 }) {
1245 state.peer_bi_streams.lock().await.remove(&result.stream_id);
1246 state
1247 .projected_stream_failures
1248 .fetch_add(1, Ordering::Relaxed);
1249 return Err(format!("send projected peer stream: {error}"));
1250 }
1251 state
1252 .projected_stream_projected
1253 .fetch_add(1, Ordering::Relaxed);
1254 Ok(ProjectedPeerBiHandoff::Handled)
1255}
1256
1257fn native_stream_trace_enabled() -> bool {
1258 cfg!(debug_assertions)
1259 || std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1")
1260}
1261
1262fn spawn_peer_bi_stream_reader(
1263 channel: Channel<Response>,
1264 stream_id: String,
1265 mut recv: PeerRecvStream,
1266) -> tokio::task::JoinHandle<()> {
1267 tokio::spawn(async move {
1268 let mut chunk = vec![0_u8; 64 * 1024];
1269 let mut close_error: Option<String> = None;
1270 loop {
1271 match recv.read(&mut chunk).await {
1272 Ok(0) => break,
1273 Ok(n) => {
1274 if native_stream_trace_enabled() {
1275 eprintln!(
1276 "[openrtc-tauri][peer-bi-stream] stream_id={} phase=plaintext-chunk bytes={}",
1277 stream_id, n
1278 );
1279 }
1280 if channel.send(Response::new(chunk[..n].to_vec())).is_err() {
1281 close_error = Some("native stream IPC channel closed".to_string());
1282 break;
1283 }
1284 }
1285 Err(error) => {
1286 if native_stream_trace_enabled() {
1287 eprintln!(
1288 "[openrtc-tauri][peer-bi-stream] stream_id={} phase=read-error error={}",
1289 stream_id, error
1290 );
1291 }
1292 close_error = Some(error.to_string());
1293 break;
1294 }
1295 }
1296 }
1297 let close = serde_json::to_string(&PeerBiStreamClosedEvent {
1298 r#type: "closed",
1299 error: close_error,
1300 })
1301 .expect("peer stream close event serializes");
1302 let _ = channel.send(Response::new(close));
1303 })
1304}
1305
1306async fn open_peer_bi_with<R: Runtime, F, Fut>(
1307 _app: tauri::AppHandle<R>,
1308 state: tauri::State<'_, OpenRtcTauriState>,
1309 peer_id: String,
1310 timeout_ms: Option<u64>,
1311 open: F,
1312) -> Result<OpenBiResult, String>
1313where
1314 F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
1315 Fut: std::future::Future<
1316 Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
1317 >,
1318{
1319 let peer_id = peer_id.trim().to_string();
1320 if peer_id.is_empty() {
1321 return Err("peerId is required".to_string());
1322 }
1323
1324 let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
1325 .await
1326 .map_err(|error| format!("open peer bi stream failed: {error}"))?;
1327 Ok(register_peer_bi_stream(&state, connection_id, remote_node_id, send, recv).await)
1328}
1329
1330#[tauri::command]
1334async fn openrtc_set_identity_credential(
1335 state: tauri::State<'_, OpenRtcTauriState>,
1336 identity_credential: Option<String>,
1337) -> Result<(), String> {
1338 state.set_identity_credential(identity_credential);
1339 if *state.presence_loop_active.lock().await {
1340 state.client().request_presence_update();
1341 }
1342 Ok(())
1343}
1344
1345fn required_app_tag(value: &str) -> Result<&str, String> {
1346 let value = value.trim();
1347 if value.len() < 5
1348 || value.len() > 80
1349 || !value.bytes().all(|byte| {
1350 byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
1351 })
1352 {
1353 return Err("OpenRTC app tag is invalid".to_string());
1354 }
1355 Ok(value)
1356}
1357
1358fn validate_public_device_jwk(value: &serde_json::Value) -> Result<(), String> {
1359 let object = value
1360 .as_object()
1361 .ok_or_else(|| "OpenRTC device public key is invalid".to_string())?;
1362 let x = object
1363 .get("x")
1364 .and_then(serde_json::Value::as_str)
1365 .unwrap_or_default();
1366 if object.get("kty").and_then(serde_json::Value::as_str) != Some("OKP")
1367 || object.get("crv").and_then(serde_json::Value::as_str) != Some("Ed25519")
1368 || object.contains_key("d")
1369 || x.len() != 43
1370 || !x
1371 .bytes()
1372 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
1373 {
1374 return Err("OpenRTC device public key must be a public Ed25519 JWK".to_string());
1375 }
1376 Ok(())
1377}
1378
1379fn required_secure_record_key(value: &str) -> Result<&str, String> {
1380 let value = value.trim();
1381 if value.is_empty()
1382 || value.len() > 512
1383 || !value
1384 .bytes()
1385 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':'))
1386 {
1387 return Err("OpenRTC secure record key is invalid".to_string());
1388 }
1389 Ok(value)
1390}
1391
1392#[tauri::command]
1393async fn openrtc_device_public_key(
1394 state: tauri::State<'_, OpenRtcTauriState>,
1395 app_tag: String,
1396) -> Result<serde_json::Value, String> {
1397 state.device_public_key(&app_tag)
1398}
1399
1400#[tauri::command]
1401async fn openrtc_sign_device_proof(
1402 state: tauri::State<'_, OpenRtcTauriState>,
1403 app_tag: String,
1404 challenge: String,
1405) -> Result<String, String> {
1406 state.sign_device_proof(&app_tag, &challenge)
1407}
1408
1409#[tauri::command]
1410async fn openrtc_delete_device_key(
1411 state: tauri::State<'_, OpenRtcTauriState>,
1412 app_tag: String,
1413) -> Result<(), String> {
1414 state.delete_device_key(&app_tag)
1415}
1416
1417#[tauri::command]
1418async fn openrtc_read_secure_record(
1419 state: tauri::State<'_, OpenRtcTauriState>,
1420 app_tag: String,
1421 key: String,
1422) -> Result<Option<String>, String> {
1423 state.read_secure_record(&app_tag, &key)
1424}
1425
1426#[tauri::command]
1427async fn openrtc_write_secure_record(
1428 state: tauri::State<'_, OpenRtcTauriState>,
1429 app_tag: String,
1430 key: String,
1431 value: String,
1432) -> Result<(), String> {
1433 state.write_secure_record(&app_tag, &key, &value)
1434}
1435
1436#[tauri::command]
1437async fn openrtc_delete_secure_record(
1438 state: tauri::State<'_, OpenRtcTauriState>,
1439 app_tag: String,
1440 key: String,
1441) -> Result<(), String> {
1442 state.delete_secure_record(&app_tag, &key)
1443}
1444
1445#[tauri::command]
1446async fn rtc_native_status(
1447 state: tauri::State<'_, OpenRtcTauriState>,
1448) -> Result<openrtc::client::RuntimeStatus, String> {
1449 Ok(state.client().runtime_status().await)
1450}
1451
1452#[tauri::command]
1453async fn get_rtc_local_device_info<R: Runtime>(
1454 app: tauri::AppHandle<R>,
1455 state: tauri::State<'_, OpenRtcTauriState>,
1456) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1457 let data_dir = app_data_dir(&app, &state)?;
1458 state
1459 .client()
1460 .init_native_device_identity(data_dir, None)
1461 .await
1462 .map_err(|error| error.to_string())
1463}
1464
1465#[tauri::command]
1466async fn update_rtc_local_device_name<R: Runtime>(
1467 app: tauri::AppHandle<R>,
1468 state: tauri::State<'_, OpenRtcTauriState>,
1469 device_name: String,
1470) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
1471 let data_dir = app_data_dir(&app, &state)?;
1472 let _ = state
1473 .client()
1474 .init_native_device_identity(data_dir, None)
1475 .await
1476 .map_err(|error| error.to_string())?;
1477 state
1478 .client()
1479 .update_native_device_name(&device_name)
1480 .await
1481 .map_err(|error| error.to_string())
1482}
1483
1484#[tauri::command]
1485async fn get_iroh_node_id(
1486 state: tauri::State<'_, OpenRtcTauriState>,
1487) -> Result<Option<String>, String> {
1488 Ok(state.client().current_node_id().await)
1489}
1490
1491#[tauri::command]
1492async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
1493 ensure_iroh_node(state.client()).await
1494}
1495
1496#[tauri::command]
1497async fn get_iroh_endpoint_ticket(
1498 state: tauri::State<'_, OpenRtcTauriState>,
1499) -> Result<String, String> {
1500 ensure_iroh_node(state.client()).await?;
1501 state
1502 .client()
1503 .endpoint_ticket()
1504 .await
1505 .map_err(|error| error.to_string())
1506}
1507
1508#[tauri::command]
1509async fn register_session_token(
1510 state: tauri::State<'_, OpenRtcTauriState>,
1511 token: String,
1512 scope: String,
1513 max_connections: u32,
1514 expires_at_ms: Option<u64>,
1515) -> Result<(), String> {
1516 if let Some(expires_at_ms) = expires_at_ms {
1517 state
1518 .client()
1519 .register_token_until(token, scope, max_connections, expires_at_ms);
1520 } else {
1521 state
1522 .client()
1523 .register_session_token(token, scope, max_connections);
1524 }
1525 Ok(())
1526}
1527
1528#[tauri::command]
1529async fn get_endpoint_ticket_with_token(
1530 state: tauri::State<'_, OpenRtcTauriState>,
1531 scope: String,
1532 max_connections: u32,
1533) -> Result<String, String> {
1534 ensure_iroh_node(state.client()).await?;
1535 state
1536 .client()
1537 .endpoint_ticket_with_token(&scope, max_connections)
1538 .await
1539 .map_err(|error| error.to_string())
1540}
1541
1542#[tauri::command]
1543async fn validate_session_token(
1544 state: tauri::State<'_, OpenRtcTauriState>,
1545 token: String,
1546 connection_id: Option<String>,
1547) -> Result<String, String> {
1548 if let Some(connection_id) = connection_id
1549 .as_deref()
1550 .map(str::trim)
1551 .filter(|value| !value.is_empty())
1552 {
1553 state
1554 .client()
1555 .validate_connection_token(&token, connection_id, None)
1556 .await
1557 .map_err(|error| error.to_string())
1558 } else {
1559 state.client().validate_session_token(&token)
1560 }
1561}
1562
1563#[tauri::command]
1564async fn revoke_session_tokens_by_scope(
1565 state: tauri::State<'_, OpenRtcTauriState>,
1566 scope: String,
1567) -> Result<Vec<String>, String> {
1568 let affected = state.client().revoke_tokens_by_scope(&scope).await;
1569 if revokes_managed_session(&scope) {
1570 if let Some(previous) = state.managed_session.lock().await.take() {
1577 previous.owner_client.stop_presence_loop();
1578 previous.owner_client.stop_auto_connect();
1579 previous.owner_client.stop_external_auto_connect().await;
1580 }
1581 *state.presence_loop_active.lock().await = false;
1582 }
1583 Ok(affected)
1584}
1585
1586#[tauri::command]
1587async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1588 state.client().stop_presence_loop();
1589 *state.presence_loop_active.lock().await = false;
1590 if let Some(active) = state.managed_session.lock().await.as_mut() {
1591 active.result.presence_started = false;
1592 }
1593 stop_subscription_by_id(&state, "internal-presence-loop").await;
1594 Ok(())
1595}
1596
1597#[tauri::command]
1603async fn register_rtc_capability(
1604 state: tauri::State<'_, OpenRtcTauriState>,
1605 capability_key: String,
1606 avenue_kind: String,
1607 avenue_id: String,
1608) -> Result<(), String> {
1609 let capability_key = required_native_capability_part(capability_key, "capability key")?;
1610 let avenue_kind = required_native_capability_part(avenue_kind, "avenue kind")?;
1611 let avenue_id = required_native_capability_part(avenue_id, "avenue id")?;
1612 if capability_key != format!("{avenue_kind}:{avenue_id}") {
1613 return Err("native capability key does not match its avenue".to_string());
1614 }
1615 state
1616 .capability_registry
1617 .lock()
1618 .await
1619 .register(capability_key, avenue_kind, avenue_id)
1620}
1621
1622#[tauri::command]
1623async fn unregister_rtc_capability(
1624 state: tauri::State<'_, OpenRtcTauriState>,
1625 capability_key: String,
1626) -> Result<(), String> {
1627 let capability_key = required_native_capability_part(capability_key, "capability key")?;
1628 let (empty, aggregate) = {
1629 let mut registry = state.capability_registry.lock().await;
1630 registry.registrations.remove(&capability_key);
1631 let empty = registry.registrations.is_empty();
1632 let aggregate = (!empty)
1633 .then(|| registry.aggregate_desired_peers())
1634 .transpose()?;
1635 (empty, aggregate)
1636 };
1637 if empty {
1638 state.client().stop_external_auto_connect().await;
1639 } else if let Some((revision, peers_json)) = aggregate {
1640 state
1641 .client()
1642 .submit_external_desired_peers(revision, &peers_json)
1643 .await
1644 .map_err(|error| error.to_string())?;
1645 }
1646 Ok(())
1647}
1648
1649#[tauri::command]
1650async fn start_rtc_managed_session<R: Runtime>(
1651 app: tauri::AppHandle<R>,
1652 state: tauri::State<'_, OpenRtcTauriState>,
1653 user_id: String,
1654 device_name: Option<String>,
1655 local_device_id: Option<String>,
1656 metadata: Option<String>,
1657 transports: Option<openrtc::client::TransportConfig>,
1658 auto_connect: Option<bool>,
1659 presence: Option<bool>,
1660) -> Result<StartSessionResult, String> {
1661 let command_started = std::time::Instant::now();
1662 let _start_guard = state.managed_session_start_guard.lock().await;
1663 let client = state.client();
1664 let app_tag = client.app_tag().to_string();
1665 eprintln!(
1666 "[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
1667 user_id,
1668 app_tag,
1669 device_name
1670 .as_deref()
1671 .map(str::trim)
1672 .map(|value| !value.is_empty())
1673 .unwrap_or(false),
1674 local_device_id
1675 .as_deref()
1676 .map(str::trim)
1677 .map(|value| !value.is_empty())
1678 .unwrap_or(false),
1679 metadata
1680 .as_deref()
1681 .map(str::trim)
1682 .map(|value| !value.is_empty())
1683 .unwrap_or(false),
1684 auto_connect.unwrap_or(false),
1685 presence.unwrap_or(false)
1686 );
1687 let data_dir = app_data_dir(&app, &state)?;
1688 let install_context = InstallContext {
1689 data_dir: data_dir.clone(),
1690 };
1691 ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
1692 .await?;
1693 if let Some(transport_config) = transports {
1694 client
1695 .update_transport_config(transport_config)
1696 .await
1697 .map_err(|error| error.to_string())?;
1698 }
1699
1700 let local_node_id = ensure_iroh_node(client.clone()).await?;
1701 let local_device = client
1702 .init_native_device_identity(data_dir, device_name.as_deref())
1703 .await
1704 .map_err(|error| error.to_string())?;
1705
1706 let effective_local_device_id =
1707 managed_session_device_id(local_device_id.as_deref(), &local_device);
1708 if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
1709 if requested != local_device.device_id {
1710 eprintln!(
1711 "[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
1712 requested, local_device.device_id
1713 );
1714 }
1715 }
1716
1717 let should_start_presence = presence.unwrap_or(false);
1718 let should_start_auto_connect = auto_connect.unwrap_or(false);
1719 if should_start_presence || should_start_auto_connect {
1720 return Err(
1721 "native Firebase coordination has been removed; publish gateway presence and submit external desired peers instead"
1722 .to_string(),
1723 );
1724 }
1725 let managed_session_key = format!("{}:{}", app_tag, effective_local_device_id);
1729 let (disposition, reused) = {
1730 let active = state.managed_session.lock().await;
1731 let disposition = managed_session_disposition(
1732 active.as_ref().map(|record| record.key.as_str()),
1733 active
1734 .as_ref()
1735 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1736 active
1737 .as_ref()
1738 .is_some_and(|record| record.result.presence_started),
1739 active
1740 .as_ref()
1741 .is_some_and(|record| record.result.auto_connect_started),
1742 &managed_session_key,
1743 should_start_presence,
1744 should_start_auto_connect,
1745 );
1746 let reused = (disposition == SessionDisposition::Reuse).then(|| {
1747 active
1748 .as_ref()
1749 .expect("managed session reuse requires an active record")
1750 .result
1751 .clone()
1752 });
1753 (disposition, reused)
1754 };
1755 if let Some(reused) = reused {
1756 eprintln!(
1757 "[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
1758 managed_session_key,
1759 command_started.elapsed().as_millis()
1760 );
1761 return Ok(reused);
1762 }
1763
1764 if disposition == SessionDisposition::Replace {
1767 let previous = state
1768 .managed_session
1769 .lock()
1770 .await
1771 .take()
1772 .expect("managed session replacement requires an active record");
1773 previous.owner_client.stop_presence_loop();
1774 previous.owner_client.stop_auto_connect();
1775 previous.owner_client.stop_external_auto_connect().await;
1776 *state.presence_loop_active.lock().await = false;
1777 }
1778 let ticket_scope = Some("user-device".to_string());
1779 let ticket = client
1780 .endpoint_ticket_with_token("user-device", 0)
1781 .await
1782 .map_err(|error| error.to_string())?;
1783
1784 let logical_local_device = local_device.clone();
1785
1786 let owner_epoch = match disposition {
1787 SessionDisposition::Refresh => {
1788 state
1789 .managed_session
1790 .lock()
1791 .await
1792 .as_ref()
1793 .expect("managed session refresh requires an active record")
1794 .owner_epoch
1795 }
1796 SessionDisposition::Start | SessionDisposition::Replace => {
1797 state.allocate_managed_session_owner_epoch()
1798 }
1799 SessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
1800 };
1801
1802 eprintln!(
1803 "[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={}",
1804 user_id,
1805 app_tag,
1806 local_node_id,
1807 logical_local_device.device_id,
1808 owner_epoch,
1809 should_start_presence,
1810 should_start_auto_connect,
1811 ticket_scope.as_deref().unwrap_or("unrestricted"),
1812 command_started.elapsed().as_millis()
1813 );
1814
1815 let result = StartSessionResult {
1816 local_node_id,
1817 ticket_scope,
1818 ticket: Some(ticket),
1819 presence_started: false,
1820 auto_connect_started: false,
1821 local_device: logical_local_device,
1822 };
1823 *state.managed_session.lock().await = Some(ManagedSessionRecord {
1824 key: managed_session_key,
1825 owner_client: client,
1826 owner_epoch,
1827 result: result.clone(),
1828 });
1829 Ok(result)
1830}
1831
1832#[tauri::command]
1833async fn start_rtc_external_auto_connect(
1834 state: tauri::State<'_, OpenRtcTauriState>,
1835 user_id: String,
1836 local_device_id: String,
1837 capability_key: Option<String>,
1838) -> Result<(), String> {
1839 if let Some(capability_key) = capability_key {
1840 let capability_key = required_native_capability_part(capability_key, "capability key")?;
1841 if !state
1842 .capability_registry
1843 .lock()
1844 .await
1845 .registrations
1846 .contains_key(&capability_key)
1847 {
1848 return Err(format!(
1849 "native capability {capability_key} is not registered"
1850 ));
1851 }
1852 }
1853 let root_user_id = format!("{}:native-root", state.client().app_tag());
1856 let _requested_capability_principal = user_id;
1857 state
1858 .client()
1859 .start_external_auto_connect(root_user_id, local_device_id)
1860 .await
1861 .map_err(|error| error.to_string())
1862}
1863
1864#[tauri::command]
1865async fn submit_rtc_desired_peers(
1866 state: tauri::State<'_, OpenRtcTauriState>,
1867 revision: u64,
1868 peers_json: String,
1869 capability_key: Option<String>,
1870) -> Result<bool, String> {
1871 let peers = serde_json::from_str::<Vec<serde_json::Value>>(&peers_json)
1872 .map_err(|error| format!("invalid capability desired-peer payload: {error}"))?;
1873 if peers.len() > 100 {
1874 return Err("capability desired-peer payload exceeds 100 peers".to_string());
1875 }
1876 let (root_revision, aggregate) = {
1877 let mut registry = state.capability_registry.lock().await;
1878 let key = match capability_key {
1879 Some(key) => required_native_capability_part(key, "capability key")?,
1880 None if registry.registrations.len() == 1 => registry
1881 .registrations
1882 .keys()
1883 .next()
1884 .cloned()
1885 .expect("one registration has one key"),
1886 None => {
1887 return Err(
1888 "capabilityKey is required when multiple native capabilities are active"
1889 .to_string(),
1890 )
1891 }
1892 };
1893 let registration = registry
1894 .registrations
1895 .get_mut(&key)
1896 .ok_or_else(|| format!("native capability {key} is not registered"))?;
1897 if revision <= registration.desired_revision {
1898 return Ok(false);
1899 }
1900 registration.desired_revision = revision;
1901 registration.desired_peers = peers;
1902 registry.aggregate_desired_peers()?
1903 };
1904 let accepted = state
1905 .client()
1906 .submit_external_desired_peers(root_revision, &aggregate)
1907 .await
1908 .map_err(|error| error.to_string())?;
1909 if accepted {
1910 replay_current_connection_states(&state).await;
1914 }
1915 Ok(accepted)
1916}
1917
1918#[tauri::command]
1919async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1920 state.client().stop_auto_connect();
1921 state.client().stop_external_auto_connect().await;
1922 if let Some(active) = state.managed_session.lock().await.as_mut() {
1923 active.result.auto_connect_started = false;
1924 }
1925 Ok(())
1926}
1927
1928#[tauri::command]
1929async fn connect_to_device(
1930 state: tauri::State<'_, OpenRtcTauriState>,
1931 device_id: Option<String>,
1932 endpoint_ticket: String,
1933 timeout_ms: Option<u64>,
1934) -> Result<openrtc::client::ManagedConnectResult, String> {
1935 let client = state.client();
1936 let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
1937 match timeout_ms {
1938 Some(timeout_ms) => {
1939 tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
1940 .await
1941 .map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
1942 .map_err(|error| error.to_string())
1943 }
1944 None => connect.await.map_err(|error| error.to_string()),
1945 }
1946}
1947
1948#[tauri::command]
1949async fn disconnect_device(
1950 state: tauri::State<'_, OpenRtcTauriState>,
1951 device_id: String,
1952 node_id_hint: Option<String>,
1953) -> Result<(), String> {
1954 state
1955 .client()
1956 .disconnect_device(&device_id, node_id_hint.as_deref())
1957 .await;
1958 Ok(())
1959}
1960
1961#[tauri::command]
1962async fn set_auto_connect_excluded(
1963 state: tauri::State<'_, OpenRtcTauriState>,
1964 device_id: String,
1965 excluded: bool,
1966) -> Result<(), String> {
1967 if excluded {
1968 state.client().exclude_peer_and_publish(&device_id).await;
1969 } else {
1970 state.client().unexclude_peer_and_publish(&device_id).await;
1971 }
1972 Ok(())
1973}
1974
1975#[tauri::command]
1976async fn set_rtc_external_auto_connect_excluded(
1977 state: tauri::State<'_, OpenRtcTauriState>,
1978 device_id: String,
1979 excluded: bool,
1980) -> Result<(), String> {
1981 state
1982 .client()
1983 .set_auto_connect_excluded(&device_id, excluded);
1984 state.client().wake_native_external_auto_connect().await;
1985 Ok(())
1986}
1987
1988#[tauri::command]
1989async fn resolve_rtc_peer_connection_records(
1990 state: tauri::State<'_, OpenRtcTauriState>,
1991 id: String,
1992) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1993 Ok(state.client().resolve_peer_connection_records(&id).await)
1994}
1995
1996#[tauri::command]
1997async fn resolve_rtc_peer_identity(
1998 state: tauri::State<'_, OpenRtcTauriState>,
1999 id: String,
2000) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
2001 Ok(state.client().peer_snapshot(&id).await)
2002}
2003
2004#[tauri::command]
2005async fn get_rtc_peer_session(
2006 state: tauri::State<'_, OpenRtcTauriState>,
2007 id: String,
2008) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2009 Ok(state.client().peer_session(&id).await)
2010}
2011
2012#[tauri::command]
2013async fn list_rtc_peer_sessions(
2014 state: tauri::State<'_, OpenRtcTauriState>,
2015) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
2016 Ok(state.client().peer_sessions().await)
2017}
2018
2019#[tauri::command]
2020async fn list_rtc_managed_connections(
2021 state: tauri::State<'_, OpenRtcTauriState>,
2022) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
2023 Ok(state.client().list_managed_connections().await)
2024}
2025
2026#[derive(Debug, Clone, Serialize)]
2027#[serde(rename_all = "camelCase")]
2028struct NotifyRtcNetworkChangeResult {
2029 retired_stale_connections: usize,
2030}
2031
2032#[tauri::command]
2033async fn notify_rtc_network_change(
2034 state: tauri::State<'_, OpenRtcTauriState>,
2035) -> Result<NotifyRtcNetworkChangeResult, String> {
2036 let retired_stale_connections = state
2037 .client()
2038 .notify_network_change()
2039 .await
2040 .map_err(|error| error.to_string())?;
2041 Ok(NotifyRtcNetworkChangeResult {
2042 retired_stale_connections,
2043 })
2044}
2045
2046#[tauri::command]
2047async fn set_rtc_transport_priority(
2048 state: tauri::State<'_, OpenRtcTauriState>,
2049 priority: Vec<openrtc::route_policy::KnownRoute>,
2050) -> Result<(), String> {
2051 let client = state.client();
2052 let mut transport_config = client.transport_config().await;
2053 transport_config.route_priority = priority;
2054 client
2055 .update_transport_config(transport_config)
2056 .await
2057 .map_err(|error| error.to_string())
2058}
2059
2060#[tauri::command]
2061async fn wait_for_rtc_settled_peer(
2062 state: tauri::State<'_, OpenRtcTauriState>,
2063 id: String,
2064 timeout_ms: Option<u64>,
2065) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
2066 Ok(state.client().wait_for_peer(&id, timeout_ms).await)
2067}
2068
2069#[tauri::command]
2070async fn list_rtc_connection_states(
2071 state: tauri::State<'_, OpenRtcTauriState>,
2072 capability_key: String,
2073) -> Result<Vec<openrtc::client::StateSnapshot>, String> {
2074 let states = state.client().connection_states().await;
2075 let capability_key = required_native_capability_part(capability_key, "capability key")?;
2076 let registry = state.capability_registry.lock().await;
2077 Ok(states
2078 .into_iter()
2079 .filter(|snapshot| {
2080 registry
2081 .capability_keys_for_state(snapshot)
2082 .iter()
2083 .any(|key| key == &capability_key)
2084 })
2085 .collect())
2086}
2087
2088#[tauri::command]
2089async fn get_rtc_connection_state(
2090 state: tauri::State<'_, OpenRtcTauriState>,
2091 connection_id: String,
2092 capability_key: String,
2093) -> Result<Option<openrtc::client::StateSnapshot>, String> {
2094 let snapshot = state.client().connection_state(&connection_id).await;
2095 let Some(snapshot) = snapshot else {
2096 return Ok(None);
2097 };
2098 let capability_key = required_native_capability_part(capability_key, "capability key")?;
2099 let included = state
2100 .capability_registry
2101 .lock()
2102 .await
2103 .capability_keys_for_state(&snapshot)
2104 .iter()
2105 .any(|key| key == &capability_key);
2106 Ok(included.then_some(snapshot))
2107}
2108
2109#[tauri::command]
2110async fn stop_rtc_subscription(
2111 state: tauri::State<'_, OpenRtcTauriState>,
2112 request_id: String,
2113) -> Result<(), String> {
2114 stop_subscription_by_id(&state, &request_id).await;
2115 Ok(())
2116}
2117
2118async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
2119 if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
2120 handle.abort();
2121 }
2122 let mut projected = state.native_projection_subscription.lock().await;
2123 if projected
2124 .as_ref()
2125 .is_some_and(|subscription| subscription.request_id == request_id)
2126 {
2127 *projected = None;
2128 }
2129}
2130
2131#[tauri::command]
2138async fn start_native_projection(
2139 state: tauri::State<'_, OpenRtcTauriState>,
2140 request_id: Option<String>,
2141 channel: Channel<NativeProjectionEvent>,
2142) -> Result<String, String> {
2143 let request_id = request_id
2144 .map(|value| value.trim().to_string())
2145 .filter(|value| !value.is_empty())
2146 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
2147 let mut subscription = state.native_projection_subscription.lock().await;
2148 if let Some(current) = subscription.as_ref() {
2149 if current.request_id == request_id {
2150 return Ok(request_id);
2151 }
2152 return Err(
2153 "OpenRTC native projection subscription already has a process owner".to_string(),
2154 );
2155 }
2156 *subscription = Some(NativeProjectionSubscription {
2157 request_id: request_id.clone(),
2158 channel,
2159 });
2160 Ok(request_id)
2161}
2162
2163#[tauri::command]
2164async fn get_projected_stream_diagnostics(
2165 state: tauri::State<'_, OpenRtcTauriState>,
2166) -> Result<ProjectedStreamDiagnostics, String> {
2167 Ok(ProjectedStreamDiagnostics {
2168 subscription_active: state.native_projection_subscription.lock().await.is_some(),
2169 offered: state.projected_stream_offered.load(Ordering::Relaxed),
2170 decoded: state.projected_stream_decoded.load(Ordering::Relaxed),
2171 projected: state.projected_stream_projected.load(Ordering::Relaxed),
2172 unhandled_non_channel: state
2173 .projected_stream_unhandled_non_channel
2174 .load(Ordering::Relaxed),
2175 unhandled_other_channel: state
2176 .projected_stream_unhandled_other_channel
2177 .load(Ordering::Relaxed),
2178 unhandled_no_subscription: state
2179 .projected_stream_unhandled_no_subscription
2180 .load(Ordering::Relaxed),
2181 failures: state.projected_stream_failures.load(Ordering::Relaxed),
2182 last_transport_stable_id: state
2183 .projected_stream_last_transport_stable_id
2184 .load(Ordering::Relaxed),
2185 validation_checks: state
2186 .projected_stream_validation_checks
2187 .load(Ordering::Relaxed),
2188 validation_rejections: state
2189 .projected_stream_validation_rejections
2190 .load(Ordering::Relaxed),
2191 last_validated_transport_stable_id: state
2192 .projected_stream_last_validated_transport_stable_id
2193 .load(Ordering::Relaxed),
2194 authorized: state.projected_stream_authorized.load(Ordering::Relaxed),
2195 unauthorized: state.projected_stream_unauthorized.load(Ordering::Relaxed),
2196 last_channel: state.projected_stream_last_channel.lock().await.clone(),
2197 last_protocol: state.projected_stream_last_protocol.lock().await.clone(),
2198 last_connection_id: state
2199 .projected_stream_last_connection_id
2200 .lock()
2201 .await
2202 .clone(),
2203 last_remote_node_id: state
2204 .projected_stream_last_remote_node_id
2205 .lock()
2206 .await
2207 .clone(),
2208 })
2209}
2210
2211#[tauri::command]
2212async fn is_current_transport_stable_id(
2213 state: tauri::State<'_, OpenRtcTauriState>,
2214 endpoint_id: String,
2215 transport_stable_id: u64,
2216) -> Result<bool, String> {
2217 state
2218 .projected_stream_validation_checks
2219 .fetch_add(1, Ordering::Relaxed);
2220 state
2221 .projected_stream_last_validated_transport_stable_id
2222 .store(transport_stable_id, Ordering::Relaxed);
2223 let is_current = state
2224 .client()
2225 .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
2226 .await
2227 .map_err(|error| error.to_string())?;
2228 if !is_current {
2229 state
2230 .projected_stream_validation_rejections
2231 .fetch_add(1, Ordering::Relaxed);
2232 }
2233 if native_stream_trace_enabled() {
2234 eprintln!(
2235 "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
2236 endpoint_id, transport_stable_id, is_current
2237 );
2238 }
2239 Ok(is_current)
2240}
2241
2242#[tauri::command]
2243async fn send_peer_message(
2244 state: tauri::State<'_, OpenRtcTauriState>,
2245 request: Request<'_>,
2246) -> Result<(), String> {
2247 let id = required_request_header(&request, PEER_ID_HEADER)?;
2248 let data = request_body_bytes(&request)?;
2249 state
2250 .client()
2251 .send_peer(&id, &data)
2252 .await
2253 .map_err(|error| error.to_string())
2254}
2255
2256#[tauri::command]
2257async fn is_peer_connected(
2258 state: tauri::State<'_, OpenRtcTauriState>,
2259 node_id: String,
2260) -> Result<bool, String> {
2261 state
2262 .client()
2263 .is_connected_str(&node_id)
2264 .await
2265 .map_err(|error| error.to_string())
2266}
2267
2268#[derive(Debug, Serialize)]
2269#[serde(rename_all = "camelCase")]
2270struct PortableMediaPublicationStart {
2271 publication_id: String,
2272 media_generation: u32,
2273 control: Vec<u8>,
2274}
2275
2276fn media_kind_name(kind: openrtc::media::MediaKind) -> &'static str {
2277 match kind {
2278 openrtc::media::MediaKind::Audio => "audio",
2279 openrtc::media::MediaKind::Video => "video",
2280 openrtc::media::MediaKind::Screen => "screen",
2281 openrtc::media::MediaKind::Data => "data",
2282 }
2283}
2284
2285fn media_codec_name(codec: openrtc::media::MediaCodec) -> &'static str {
2286 match codec {
2287 openrtc::media::MediaCodec::Opus => "opus",
2288 openrtc::media::MediaCodec::H264 => "h264",
2289 openrtc::media::MediaCodec::Vp8 => "vp8",
2290 openrtc::media::MediaCodec::Vp9 => "vp9",
2291 openrtc::media::MediaCodec::Av1 => "av1",
2292 openrtc::media::MediaCodec::Pcm => "pcm",
2293 openrtc::media::MediaCodec::Opaque => "opaque",
2294 }
2295}
2296
2297fn parse_media_publication_id(value: &str) -> Result<openrtc::media::PublicationId, String> {
2298 value
2299 .trim()
2300 .parse()
2301 .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())
2302}
2303
2304#[tauri::command]
2305#[allow(clippy::too_many_arguments)]
2306async fn begin_openrtc_media_publication(
2307 state: tauri::State<'_, OpenRtcTauriState>,
2308 publication_id: Option<String>,
2309 kind: String,
2310 codec: String,
2311 clock_rate: u32,
2312 coded_width: Option<u32>,
2313 coded_height: Option<u32>,
2314 channels: Option<u16>,
2315) -> Result<PortableMediaPublicationStart, String> {
2316 let publication_id = match publication_id.as_deref().map(str::trim) {
2317 Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
2318 _ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
2319 };
2320 let kind: openrtc::media::MediaKind = kind
2321 .parse()
2322 .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
2323 if !matches!(
2324 kind,
2325 openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
2326 ) {
2327 return Err("browser media publications must be audio or video".to_string());
2328 }
2329 let codec = codec
2330 .parse()
2331 .map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
2332 let publication = openrtc::media::MediaPublicationConfig {
2333 publication_id,
2334 media_generation: 0,
2335 kind,
2336 codec,
2337 clock_rate,
2338 coded_width,
2339 coded_height,
2340 channels,
2341 };
2342 let (publication, control) = state
2343 .portable_media
2344 .lock()
2345 .await
2346 .begin_publication(publication)
2347 .map_err(|error| error.to_string())?;
2348 Ok(PortableMediaPublicationStart {
2349 publication_id: publication.publication_id.to_string(),
2350 media_generation: publication.media_generation,
2351 control,
2352 })
2353}
2354
2355#[tauri::command]
2356#[allow(clippy::too_many_arguments)]
2357async fn encode_openrtc_media_sample(
2358 state: tauri::State<'_, OpenRtcTauriState>,
2359 request: Request<'_>,
2360) -> Result<Response, String> {
2361 let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
2362 let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
2363 openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
2364 .map_err(|error| error.to_string())?;
2365 let duration_us = parsed_request_header(&request, MEDIA_DURATION_US_HEADER)?;
2366 let keyframe = parsed_request_header(&request, MEDIA_KEYFRAME_HEADER)?;
2367 let discardable = parsed_request_header(&request, MEDIA_DISCARDABLE_HEADER)?;
2368 let payload = request_body_bytes(&request)?;
2369 let encoded = state
2370 .portable_media
2371 .lock()
2372 .await
2373 .encode_sample(
2374 parse_media_publication_id(&publication_id)?,
2375 openrtc::media::EncodedMediaSample {
2376 timestamp_us,
2377 duration_us,
2378 keyframe,
2379 discardable,
2380 payload,
2381 },
2382 )
2383 .map_err(|error| error.to_string())?;
2384 Ok(Response::new(encoded))
2385}
2386
2387#[tauri::command]
2388async fn pause_openrtc_media_publication(
2389 state: tauri::State<'_, OpenRtcTauriState>,
2390 publication_id: String,
2391) -> Result<(), String> {
2392 state
2393 .portable_media
2394 .lock()
2395 .await
2396 .pause_publication(parse_media_publication_id(&publication_id)?)
2397 .map_err(|error| error.to_string())
2398}
2399
2400#[tauri::command]
2401async fn retire_openrtc_media_publication(
2402 state: tauri::State<'_, OpenRtcTauriState>,
2403 publication_id: String,
2404) -> Result<(), String> {
2405 state
2406 .portable_media
2407 .lock()
2408 .await
2409 .retire_publication(parse_media_publication_id(&publication_id)?);
2410 Ok(())
2411}
2412
2413#[tauri::command]
2414async fn retire_openrtc_media_receiver(
2415 state: tauri::State<'_, OpenRtcTauriState>,
2416 publication_id: String,
2417 media_generation: u32,
2418) -> Result<(), String> {
2419 state.portable_media.lock().await.retire_receiver(
2420 parse_media_publication_id(&publication_id)?,
2421 media_generation,
2422 );
2423 Ok(())
2424}
2425
2426#[tauri::command]
2427async fn decode_openrtc_media_chunk(
2428 state: tauri::State<'_, OpenRtcTauriState>,
2429 request: Request<'_>,
2430) -> Result<Response, String> {
2431 let encoded = request_body_bytes(&request)?;
2432 let chunk = state
2433 .portable_media
2434 .lock()
2435 .await
2436 .decode_chunk(&encoded)
2437 .map_err(|error| error.to_string())?;
2438 openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
2439 .and_then(|_| {
2440 openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
2441 })
2442 .map_err(|error| error.to_string())?;
2443 let metadata = serde_json::to_vec(&serde_json::json!({
2444 "publicationId": chunk.publication_id.to_string(),
2445 "mediaGeneration": chunk.media_generation,
2446 "sequence": chunk.sequence,
2447 "timestampUs": chunk.timestamp_us,
2448 "durationUs": chunk.duration_us,
2449 "kind": media_kind_name(chunk.kind),
2450 "codec": media_codec_name(chunk.codec),
2451 "keyframe": chunk.keyframe,
2452 "discardable": chunk.discardable,
2453 }))
2454 .map_err(|error| format!("encode media IPC metadata failed: {error}"))?;
2455 let metadata_len =
2456 u32::try_from(metadata.len()).map_err(|_| "media IPC metadata is too large".to_string())?;
2457 let mut response = Vec::with_capacity(4 + metadata.len() + chunk.payload.len());
2458 response.extend_from_slice(&metadata_len.to_be_bytes());
2459 response.extend_from_slice(&metadata);
2460 response.extend_from_slice(&chunk.payload);
2461 Ok(Response::new(response))
2462}
2463
2464#[tauri::command]
2465async fn decode_openrtc_media_control(
2466 state: tauri::State<'_, OpenRtcTauriState>,
2467 request: Request<'_>,
2468) -> Result<serde_json::Value, String> {
2469 let encoded = request_body_bytes(&request)?;
2470 let control = state
2471 .portable_media
2472 .lock()
2473 .await
2474 .decode_control(&encoded)
2475 .map_err(|error| error.to_string())?;
2476 Ok(match control {
2477 openrtc::media::MediaControlFrame::Publish {
2478 publication_id,
2479 media_generation,
2480 kind,
2481 codec,
2482 clock_rate,
2483 coded_width,
2484 coded_height,
2485 channels,
2486 } => serde_json::json!({
2487 "type": "publish",
2488 "publicationId": publication_id.to_string(),
2489 "mediaGeneration": media_generation,
2490 "kind": media_kind_name(kind),
2491 "codec": media_codec_name(codec),
2492 "clockRate": clock_rate,
2493 "codedWidth": coded_width,
2494 "codedHeight": coded_height,
2495 "channels": channels,
2496 }),
2497 openrtc::media::MediaControlFrame::SetEnabled {
2498 publication_id,
2499 media_generation,
2500 enabled,
2501 } => serde_json::json!({
2502 "type": "set-enabled",
2503 "publicationId": publication_id.to_string(),
2504 "mediaGeneration": media_generation,
2505 "enabled": enabled,
2506 }),
2507 openrtc::media::MediaControlFrame::RequestKeyframe {
2508 publication_id,
2509 media_generation,
2510 } => serde_json::json!({
2511 "type": "request-keyframe",
2512 "publicationId": publication_id.to_string(),
2513 "mediaGeneration": media_generation,
2514 }),
2515 openrtc::media::MediaControlFrame::Stop {
2516 publication_id,
2517 media_generation,
2518 reason,
2519 } => serde_json::json!({
2520 "type": "stop",
2521 "publicationId": publication_id.to_string(),
2522 "mediaGeneration": media_generation,
2523 "reason": reason,
2524 }),
2525 })
2526}
2527
2528#[tauri::command]
2529async fn open_peer_bi_stream<R: Runtime>(
2530 app: tauri::AppHandle<R>,
2531 state: tauri::State<'_, OpenRtcTauriState>,
2532 peer_id: String,
2533 timeout_ms: Option<u64>,
2534) -> Result<OpenBiResult, String> {
2535 open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
2536 client.open_peer_bi(&peer_id, timeout_ms).await
2537 })
2538 .await
2539}
2540
2541#[tauri::command]
2542async fn open_peer_bi_transport_only_stream<R: Runtime>(
2543 app: tauri::AppHandle<R>,
2544 state: tauri::State<'_, OpenRtcTauriState>,
2545 peer_id: String,
2546 timeout_ms: Option<u64>,
2547) -> Result<OpenBiResult, String> {
2548 open_peer_bi_with(
2549 app,
2550 state,
2551 peer_id,
2552 timeout_ms,
2553 |client, peer_id, timeout_ms| async move {
2554 client
2555 .open_peer_bi_transport_only(&peer_id, timeout_ms)
2556 .await
2557 .map(|(connection_id, remote_node_id, send, recv)| {
2558 (
2559 connection_id,
2560 remote_node_id,
2561 PeerSendStream::plain(send),
2562 PeerRecvStream::plain(recv),
2563 )
2564 })
2565 },
2566 )
2567 .await
2568}
2569
2570#[tauri::command]
2571async fn open_peer_uni_stream(
2572 state: tauri::State<'_, OpenRtcTauriState>,
2573 peer_id: String,
2574 timeout_ms: Option<u64>,
2575) -> Result<OpenUniResult, String> {
2576 let peer_id = peer_id.trim().to_string();
2577 if peer_id.is_empty() {
2578 return Err("peerId is required".to_string());
2579 }
2580 let (connection_id, remote_node_id, send) = state
2581 .client()
2582 .open_peer_uni(&peer_id, timeout_ms)
2583 .await
2584 .map_err(|error| format!("open peer uni stream failed: {error}"))?;
2585 let stream_id = uuid::Uuid::new_v4().to_string();
2586 state
2587 .peer_uni_streams
2588 .lock()
2589 .await
2590 .insert(stream_id.clone(), Arc::new(Mutex::new(Some(send))));
2591 Ok(OpenUniResult {
2592 stream_id,
2593 connection_id,
2594 remote_node_id,
2595 })
2596}
2597
2598#[tauri::command]
2599async fn write_peer_bi_stream(
2600 state: tauri::State<'_, OpenRtcTauriState>,
2601 request: Request<'_>,
2602) -> Result<(), String> {
2603 let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
2604 let bytes = request_body_bytes(&request)?;
2605
2606 let send = {
2607 let streams = state.peer_bi_streams.lock().await;
2608 streams
2609 .get(&stream_id)
2610 .map(|handle| handle.send.clone())
2611 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
2612 };
2613 let result = send
2614 .lock()
2615 .await
2616 .as_mut()
2617 .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
2618 .write_all(&bytes)
2619 .await
2620 .map_err(|error| format!("write peer bi stream failed: {error}"));
2621 result
2622}
2623
2624#[tauri::command]
2625async fn start_peer_bi_stream_read<R: Runtime>(
2626 _app: tauri::AppHandle<R>,
2627 state: tauri::State<'_, OpenRtcTauriState>,
2628 stream_id: String,
2629 channel: Channel<Response>,
2630) -> Result<(), String> {
2631 let stream_id = stream_id.trim().to_string();
2632 if stream_id.is_empty() {
2633 return Err("streamId is required".to_string());
2634 }
2635
2636 let mut streams = state.peer_bi_streams.lock().await;
2637 let handle = streams
2638 .get_mut(&stream_id)
2639 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
2640 if handle.read_task.is_some() {
2641 return Ok(());
2642 }
2643 let recv = handle
2644 .recv
2645 .take()
2646 .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
2647 handle.read_task = Some(spawn_peer_bi_stream_reader(channel, stream_id, recv));
2648 Ok(())
2649}
2650
2651#[tauri::command]
2657async fn cancel_peer_bi_stream_read(
2658 state: tauri::State<'_, OpenRtcTauriState>,
2659 stream_id: String,
2660) -> Result<(), String> {
2661 let stream_id = stream_id.trim().to_string();
2662 if stream_id.is_empty() {
2663 return Err("streamId is required".to_string());
2664 }
2665
2666 let mut streams = state.peer_bi_streams.lock().await;
2667 let handle = streams
2668 .get_mut(&stream_id)
2669 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
2670 handle.recv.take();
2671 if let Some(read_task) = handle.read_task.take() {
2672 read_task.abort();
2673 }
2674 Ok(())
2675}
2676
2677#[tauri::command]
2678async fn close_peer_bi_stream(
2679 state: tauri::State<'_, OpenRtcTauriState>,
2680 stream_id: String,
2681) -> Result<(), String> {
2682 let stream_id = stream_id.trim().to_string();
2683 if stream_id.is_empty() {
2684 return Err("streamId is required".to_string());
2685 }
2686
2687 let handle = {
2688 let mut streams = state.peer_bi_streams.lock().await;
2689 streams.remove(&stream_id)
2690 };
2691 if let Some(handle) = handle {
2692 let mut send = handle.send.lock().await;
2695 if let Some(send) = send.take() {
2696 let _ = send.finish();
2697 }
2698 drop(send);
2699 if let Some(read_task) = handle.read_task {
2700 read_task.abort();
2701 }
2702 }
2703 Ok(())
2704}
2705
2706#[tauri::command]
2707async fn write_peer_uni_stream(
2708 state: tauri::State<'_, OpenRtcTauriState>,
2709 request: Request<'_>,
2710) -> Result<(), String> {
2711 let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
2712 let bytes = request_body_bytes(&request)?;
2713 let send = {
2714 let streams = state.peer_uni_streams.lock().await;
2715 streams
2716 .get(&stream_id)
2717 .cloned()
2718 .ok_or_else(|| format!("peer uni stream not found: {stream_id}"))?
2719 };
2720 let result = send
2721 .lock()
2722 .await
2723 .as_mut()
2724 .ok_or_else(|| format!("peer uni stream already closed: {stream_id}"))?
2725 .write_all(&bytes)
2726 .await
2727 .map_err(|error| format!("write peer uni stream failed: {error}"));
2728 result
2729}
2730
2731#[tauri::command]
2732async fn close_peer_uni_stream(
2733 state: tauri::State<'_, OpenRtcTauriState>,
2734 stream_id: String,
2735) -> Result<(), String> {
2736 let stream_id = stream_id.trim().to_string();
2737 if stream_id.is_empty() {
2738 return Err("streamId is required".to_string());
2739 }
2740 let send = state.peer_uni_streams.lock().await.remove(&stream_id);
2741 if let Some(send) = send {
2742 if let Some(send) = send.lock().await.take() {
2743 let _ = send.finish();
2744 }
2745 }
2746 Ok(())
2747}
2748
2749#[cfg(test)]
2750mod tests {
2751 use super::*;
2752 use std::collections::HashMap as StdHashMap;
2753 use std::sync::atomic::{AtomicUsize, Ordering};
2754 use std::sync::Mutex as StdMutex;
2755 use tokio::sync::oneshot;
2756
2757 struct FakeTransportInstaller {
2758 installs: Arc<AtomicUsize>,
2759 }
2760
2761 #[derive(Default)]
2762 struct FakeDeviceKeySigner {
2763 records: StdMutex<StdHashMap<String, String>>,
2764 }
2765
2766 impl DeviceKeySigner for FakeDeviceKeySigner {
2767 fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
2768 Ok(serde_json::json!({
2769 "kty": "OKP",
2770 "crv": "Ed25519",
2771 "x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
2772 }))
2773 }
2774
2775 fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
2776 assert_eq!(challenge, b"openrtc:v2:test");
2777 Ok(vec![7; 64])
2778 }
2779
2780 fn delete(&self, _app_tag: &str) -> Result<(), String> {
2781 Ok(())
2782 }
2783
2784 fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
2785 Ok(self
2786 .records
2787 .lock()
2788 .unwrap()
2789 .get(&format!("{app_tag}:{key}"))
2790 .cloned())
2791 }
2792
2793 fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
2794 self.records
2795 .lock()
2796 .unwrap()
2797 .insert(format!("{app_tag}:{key}"), value.to_string());
2798 Ok(())
2799 }
2800
2801 fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
2802 self.records
2803 .lock()
2804 .unwrap()
2805 .remove(&format!("{app_tag}:{key}"));
2806 Ok(())
2807 }
2808 }
2809
2810 impl TransportInstaller for FakeTransportInstaller {
2811 fn id(&self) -> &'static str {
2812 "ble"
2813 }
2814
2815 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2816 config.ble.as_ref().is_some_and(|ble| ble.enabled)
2817 }
2818
2819 fn install(
2820 &self,
2821 _client: Arc<openrtc::client::Client>,
2822 _config: openrtc::client::TransportConfig,
2823 _context: InstallContext,
2824 ) -> InstallFuture {
2825 self.installs.fetch_add(1, Ordering::SeqCst);
2826 Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
2827 }
2828 }
2829
2830 #[derive(Debug)]
2831 struct UnavailableTransportInstaller;
2832
2833 impl TransportInstaller for UnavailableTransportInstaller {
2834 fn id(&self) -> &'static str {
2835 "ble"
2836 }
2837
2838 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
2839 config.ble.as_ref().is_some_and(|ble| ble.enabled)
2840 }
2841
2842 fn install(
2843 &self,
2844 _client: Arc<openrtc::client::Client>,
2845 _config: openrtc::client::TransportConfig,
2846 _context: InstallContext,
2847 ) -> InstallFuture {
2848 Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
2849 }
2850 }
2851
2852 fn requested_ble_config() -> openrtc::client::TransportConfig {
2853 openrtc::client::TransportConfig {
2854 ble: Some(openrtc::client::BleConfig {
2855 enabled: true,
2856 ..openrtc::client::BleConfig::default()
2857 }),
2858 ..openrtc::client::TransportConfig::default()
2859 }
2860 }
2861
2862 fn test_config() -> OpenRtcTauriConfig {
2863 OpenRtcTauriConfig {
2864 api_key: format!("pk_test_{}", "a".repeat(40)),
2865 ..OpenRtcTauriConfig::default()
2866 }
2867 }
2868
2869 #[test]
2870 fn native_capability_registry_isolates_duplicate_peer_channels() {
2871 let mut registry = NativeCapabilityRegistry::default();
2872 registry
2873 .register(
2874 "devices:user-1".to_string(),
2875 "devices".to_string(),
2876 "user-1".to_string(),
2877 )
2878 .expect("devices registration");
2879 registry
2880 .register(
2881 "space:room-1".to_string(),
2882 "space".to_string(),
2883 "room-1".to_string(),
2884 )
2885 .expect("space registration");
2886 registry
2887 .registrations
2888 .get_mut("devices:user-1")
2889 .unwrap()
2890 .desired_peers = vec![serde_json::json!({
2891 "deviceId": "shared-peer",
2892 "nodeId": "device-route"
2893 })];
2894 registry
2895 .registrations
2896 .get_mut("space:room-1")
2897 .unwrap()
2898 .desired_peers = vec![serde_json::json!({
2899 "deviceId": "shared-peer",
2900 "nodeId": "space-route"
2901 })];
2902
2903 assert_eq!(
2904 registry.capability_keys_for_identity(
2905 Some("shared-peer"),
2906 Some("shared-peer"),
2907 None,
2908 None,
2909 ),
2910 vec!["devices:user-1".to_string(), "space:room-1".to_string()]
2911 );
2912
2913 let explicit_channel = openrtc::stream_metadata::ChannelMetadata {
2914 channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
2915 metadata: Some(serde_json::Map::from_iter([(
2916 "openrtcCapability".to_string(),
2917 serde_json::json!("space:room-1"),
2918 )])),
2919 };
2920 assert_eq!(
2921 registry.projected_stream_capability(&explicit_channel),
2922 Some("space:room-1".to_string())
2923 );
2924
2925 let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
2926 channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
2927 metadata: None,
2928 };
2929 assert_eq!(
2930 registry.projected_stream_capability(&unscoped_channel),
2931 None
2932 );
2933 }
2934
2935 #[test]
2936 fn native_capability_registry_aggregates_without_duplicating_root_peers() {
2937 let mut registry = NativeCapabilityRegistry::default();
2938 for (key, kind, id) in [
2939 ("devices:user-1", "devices", "user-1"),
2940 ("space:room-1", "space", "room-1"),
2941 ] {
2942 registry
2943 .register(key.to_string(), kind.to_string(), id.to_string())
2944 .expect("capability registration");
2945 registry.registrations.get_mut(key).unwrap().desired_peers =
2946 vec![serde_json::json!({"deviceId": "shared-peer", "ticket": "ticket-a"})];
2947 }
2948
2949 let (first_revision, peers_json) = registry.aggregate_desired_peers().unwrap();
2950 let peers: Vec<serde_json::Value> = serde_json::from_str(&peers_json).unwrap();
2951 assert_eq!(first_revision, 1);
2952 assert_eq!(peers.len(), 1, "one physical peer must be dialed once");
2953
2954 registry
2955 .register(
2956 "devices:user-1".to_string(),
2957 "devices".to_string(),
2958 "user-1".to_string(),
2959 )
2960 .expect("exact registration is idempotent");
2961 assert!(registry
2962 .register(
2963 "devices:user-1".to_string(),
2964 "space".to_string(),
2965 "room-1".to_string(),
2966 )
2967 .is_err());
2968 }
2969
2970 #[test]
2971 fn native_capability_registry_scopes_peer_data_and_rejects_unscoped_streams() {
2972 let mut registry = NativeCapabilityRegistry::default();
2973 registry
2974 .register(
2975 "devices:user-1".to_string(),
2976 "devices".to_string(),
2977 "user-1".to_string(),
2978 )
2979 .unwrap();
2980 registry
2981 .registrations
2982 .get_mut("devices:user-1")
2983 .unwrap()
2984 .desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
2985
2986 let explicit = openrtc::client::NativePeerDataEvent {
2987 connection_id: "peer-1".to_string(),
2988 remote_node_id: None,
2989 transport: "webrtc".to_string(),
2990 transport_stable_id: 1,
2991 transport_generation: 1,
2992 route_generation: 1,
2993 payload: serde_json::to_vec(&serde_json::json!({
2994 "capability": "devices:user-1",
2995 "kind": "raw",
2996 "payload": [1, 2, 3]
2997 }))
2998 .unwrap(),
2999 };
3000 assert_eq!(
3001 registry.capability_keys_for_peer_data(&explicit),
3002 vec!["devices:user-1".to_string()]
3003 );
3004
3005 let unknown = openrtc::client::NativePeerDataEvent {
3006 payload: serde_json::to_vec(&serde_json::json!({
3007 "capability": "space:unknown",
3008 "kind": "raw"
3009 }))
3010 .unwrap(),
3011 ..explicit.clone()
3012 };
3013 assert!(registry.capability_keys_for_peer_data(&unknown).is_empty());
3014
3015 let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
3016 channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
3017 metadata: None,
3018 };
3019 assert_eq!(
3020 registry.projected_stream_capability(&unscoped_channel),
3021 None,
3022 "stream ownership must never be inferred from capability count"
3023 );
3024 }
3025
3026 #[test]
3027 fn native_device_proof_uses_host_signer_without_exporting_private_key() {
3028 let state = OpenRtcTauriState::new(test_config())
3029 .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
3030 let public = state
3031 .device_public_key("app_native_test")
3032 .expect("public key");
3033 assert_eq!(public["kty"], "OKP");
3034 assert_eq!(public["crv"], "Ed25519");
3035 assert!(public.get("d").is_none());
3036 let signature = state
3037 .sign_device_proof("app_native_test", "openrtc:v2:test")
3038 .expect("signature");
3039 assert_eq!(
3040 base64::engine::general_purpose::URL_SAFE_NO_PAD
3041 .decode(signature)
3042 .expect("base64"),
3043 vec![7; 64]
3044 );
3045 }
3046
3047 #[test]
3048 fn native_certificate_records_round_trip_through_host_secure_storage() {
3049 let state = OpenRtcTauriState::new(test_config())
3050 .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
3051 let app_tag = "app_native_test";
3052 let key = "openrtc:v2:device-session:abc:device-1";
3053 assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
3054 state
3055 .write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
3056 .unwrap();
3057 assert_eq!(
3058 state.read_secure_record(app_tag, key).unwrap().as_deref(),
3059 Some("{\"token\":\"bound\"}")
3060 );
3061 state.delete_secure_record(app_tag, key).unwrap();
3062 assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
3063 }
3064
3065 #[test]
3066 fn native_device_proof_fails_closed_without_secure_host_signer() {
3067 let state = OpenRtcTauriState::new(test_config());
3068 assert!(state
3069 .device_public_key("app_native_test")
3070 .expect_err("missing signer must fail")
3071 .contains("secure-storage signer"));
3072 }
3073
3074 fn native_transport_install_context() -> InstallContext {
3075 InstallContext {
3076 data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
3077 }
3078 }
3079
3080 fn managed_session_test_result(device_id: &str) -> StartSessionResult {
3081 StartSessionResult {
3082 local_node_id: format!("node-{device_id}"),
3083 ticket_scope: Some("user-device".to_string()),
3084 ticket: None,
3085 presence_started: true,
3086 auto_connect_started: true,
3087 local_device: openrtc::native_device::NativeDeviceIdentity {
3088 device_id: device_id.to_string(),
3089 device_name: "Test Device".to_string(),
3090 created_at_ms: 1,
3091 updated_at_ms: 1,
3092 name_source: None,
3093 system_info: None,
3094 },
3095 }
3096 }
3097
3098 async fn simulate_managed_session_start(
3099 state: Arc<OpenRtcTauriState>,
3100 device_id: String,
3101 pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
3102 starts: Arc<AtomicUsize>,
3103 stopped_owners: Arc<Mutex<Vec<usize>>>,
3104 ) -> SessionDisposition {
3105 let _start_guard = state.managed_session_start_guard.lock().await;
3106 let client = state.client();
3107 let key = format!("{}:{device_id}", client.app_tag());
3108 let (disposition, active_epoch) = {
3109 let active = state.managed_session.lock().await;
3110 let disposition = managed_session_disposition(
3111 active.as_ref().map(|record| record.key.as_str()),
3112 active
3113 .as_ref()
3114 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
3115 active
3116 .as_ref()
3117 .is_some_and(|record| record.result.presence_started),
3118 active
3119 .as_ref()
3120 .is_some_and(|record| record.result.auto_connect_started),
3121 &key,
3122 true,
3123 true,
3124 );
3125 (
3126 disposition,
3127 active.as_ref().map(|record| record.owner_epoch),
3128 )
3129 };
3130
3131 if disposition == SessionDisposition::Reuse {
3132 return disposition;
3133 }
3134
3135 if disposition == SessionDisposition::Replace {
3136 let previous = state
3137 .managed_session
3138 .lock()
3139 .await
3140 .take()
3141 .expect("replacement must have a previous owner");
3142 stopped_owners
3143 .lock()
3144 .await
3145 .push(Arc::as_ptr(&previous.owner_client) as usize);
3146 }
3147
3148 starts.fetch_add(1, Ordering::SeqCst);
3149 if let Some((entered, release)) = pause {
3150 entered.send(()).expect("start observer must be waiting");
3151 release.await.expect("start release must be sent");
3152 }
3153
3154 let owner_epoch = match disposition {
3155 SessionDisposition::Refresh => {
3156 active_epoch.expect("refresh must preserve the active owner epoch")
3157 }
3158 SessionDisposition::Start | SessionDisposition::Replace => {
3159 state.allocate_managed_session_owner_epoch()
3160 }
3161 SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
3162 };
3163 let result = managed_session_test_result(&device_id);
3164 *state.managed_session.lock().await = Some(ManagedSessionRecord {
3165 key,
3166 owner_client: client,
3167 owner_epoch,
3168 result,
3169 });
3170 disposition
3171 }
3172
3173 #[tokio::test]
3174 async fn concurrent_same_key_start_reuses_one_owner_epoch() {
3175 let state = Arc::new(OpenRtcTauriState::new(test_config()));
3176 let starts = Arc::new(AtomicUsize::new(0));
3177 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
3178 let (first_entered_tx, first_entered_rx) = oneshot::channel();
3179 let (first_release_tx, first_release_rx) = oneshot::channel();
3180
3181 let first = tokio::spawn(simulate_managed_session_start(
3182 state.clone(),
3183 "same-device".to_string(),
3184 Some((first_entered_tx, first_release_rx)),
3185 starts.clone(),
3186 stopped_owners.clone(),
3187 ));
3188 first_entered_rx.await.expect("first start must pause");
3189
3190 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
3191 let second_state = state.clone();
3192 let second_starts = starts.clone();
3193 let second_stopped_owners = stopped_owners.clone();
3194 let second = tokio::spawn(async move {
3195 second_attempting_tx.send(()).unwrap();
3196 simulate_managed_session_start(
3197 second_state,
3198 "same-device".to_string(),
3199 None,
3200 second_starts,
3201 second_stopped_owners,
3202 )
3203 .await
3204 });
3205 second_attempting_rx.await.unwrap();
3206 tokio::task::yield_now().await;
3207 assert_eq!(starts.load(Ordering::SeqCst), 1);
3208
3209 first_release_tx.send(()).unwrap();
3210 assert_eq!(
3211 first.await.expect("first start task"),
3212 SessionDisposition::Start
3213 );
3214 assert_eq!(
3215 second.await.expect("second start task"),
3216 SessionDisposition::Reuse
3217 );
3218 assert_eq!(starts.load(Ordering::SeqCst), 1);
3219
3220 let active = state
3221 .managed_session
3222 .lock()
3223 .await
3224 .clone()
3225 .expect("active owner");
3226 assert_eq!(active.owner_epoch, 1);
3227 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
3228 assert!(stopped_owners.lock().await.is_empty());
3229 }
3230
3231 #[tokio::test]
3232 async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
3233 let state = Arc::new(OpenRtcTauriState::new(test_config()));
3234 let starts = Arc::new(AtomicUsize::new(0));
3235 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
3236 let (first_entered_tx, first_entered_rx) = oneshot::channel();
3237 let (first_release_tx, first_release_rx) = oneshot::channel();
3238
3239 let first = tokio::spawn(simulate_managed_session_start(
3240 state.clone(),
3241 "first-device".to_string(),
3242 Some((first_entered_tx, first_release_rx)),
3243 starts.clone(),
3244 stopped_owners.clone(),
3245 ));
3246 first_entered_rx.await.expect("first start must pause");
3247 let first_owner = state.client();
3248
3249 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
3250 let second_state = state.clone();
3251 let second_starts = starts.clone();
3252 let second_stopped_owners = stopped_owners.clone();
3253 let second = tokio::spawn(async move {
3254 second_attempting_tx.send(()).unwrap();
3255 simulate_managed_session_start(
3256 second_state,
3257 "second-device".to_string(),
3258 None,
3259 second_starts,
3260 second_stopped_owners,
3261 )
3262 .await
3263 });
3264 second_attempting_rx.await.unwrap();
3265 tokio::task::yield_now().await;
3266 assert!(Arc::ptr_eq(&first_owner, &state.client()));
3267
3268 first_release_tx.send(()).unwrap();
3269 assert_eq!(
3270 first.await.expect("first start task"),
3271 SessionDisposition::Start
3272 );
3273 assert_eq!(
3274 second.await.expect("replacement start task"),
3275 SessionDisposition::Replace
3276 );
3277 assert_eq!(starts.load(Ordering::SeqCst), 2);
3278
3279 let active = state
3280 .managed_session
3281 .lock()
3282 .await
3283 .clone()
3284 .expect("active owner");
3285 assert_eq!(active.owner_epoch, 2);
3286 assert!(active.key.contains("second-device"));
3287 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
3288 assert_eq!(
3289 stopped_owners.lock().await.as_slice(),
3290 [Arc::as_ptr(&first_owner) as usize]
3291 );
3292 assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
3293 }
3294
3295 #[test]
3296 fn derives_app_tag_from_api_key_by_default() {
3297 let config = OpenRtcTauriConfig {
3298 api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
3299 ..OpenRtcTauriConfig::default()
3300 };
3301
3302 assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
3303 }
3304
3305 #[tokio::test]
3306 async fn requested_native_transport_without_installer_keeps_base_route_available() {
3307 let state = OpenRtcTauriState::new(test_config());
3308 let client = state.client();
3309 ensure_requested_native_transports(
3310 &state,
3311 &client,
3312 Some(&requested_ble_config()),
3313 &native_transport_install_context(),
3314 )
3315 .await
3316 .expect("an unavailable optional transport must not fail base Iroh startup");
3317
3318 assert!(state.installed_native_transports.lock().await.is_empty());
3319 }
3320
3321 #[tokio::test]
3322 async fn native_transport_installer_is_idempotent_for_one_client() {
3323 let installs = Arc::new(AtomicUsize::new(0));
3324 let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
3325 Arc::new(FakeTransportInstaller {
3326 installs: installs.clone(),
3327 }),
3328 );
3329 let client = state.client();
3330 let config = requested_ble_config();
3331
3332 ensure_requested_native_transports(
3333 &state,
3334 &client,
3335 Some(&config),
3336 &native_transport_install_context(),
3337 )
3338 .await
3339 .expect("first install");
3340 ensure_requested_native_transports(
3341 &state,
3342 &client,
3343 Some(&config),
3344 &native_transport_install_context(),
3345 )
3346 .await
3347 .expect("idempotent install");
3348
3349 assert_eq!(installs.load(Ordering::SeqCst), 1);
3350 }
3351
3352 #[tokio::test]
3353 async fn unavailable_native_transport_keeps_base_route_available() {
3354 let state = OpenRtcTauriState::new(test_config())
3355 .with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
3356 let client = state.client();
3357
3358 ensure_requested_native_transports(
3359 &state,
3360 &client,
3361 Some(&requested_ble_config()),
3362 &native_transport_install_context(),
3363 )
3364 .await
3365 .expect("optional transport installation failure must not fail base Iroh startup");
3366
3367 assert!(state.installed_native_transports.lock().await.is_empty());
3368 }
3369
3370 #[test]
3371 fn invalid_or_missing_api_key_fails_closed() {
3372 assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
3373 let mut config = OpenRtcTauriConfig::default();
3374 config.api_key = "pk_test_short".to_string();
3375 assert!(config.validated_api_key().is_err());
3376 }
3377
3378 #[test]
3379 fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
3380 let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
3381 let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
3382 let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
3383 let api_key = format!("pk_test_{}", "c".repeat(40));
3384 std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
3385 std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
3386 std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
3387
3388 let config = OpenRtcTauriConfig::from_env();
3389
3390 assert_eq!(config.api_key, api_key);
3391 assert_eq!(
3392 config.app_tag().unwrap(),
3393 openrtc::app_tag_from_api_key(&api_key)
3394 );
3395
3396 match previous_api_key {
3397 Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
3398 None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
3399 }
3400 match previous_project {
3401 Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
3402 None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
3403 }
3404 match previous_app_tag {
3405 Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
3406 None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
3407 }
3408 }
3409
3410 #[test]
3411 fn native_state_has_one_constructor_owned_app_identity() {
3412 let config = test_config();
3413 let expected = openrtc::app_tag_from_api_key(&config.api_key);
3414 let state = OpenRtcTauriState::new(config);
3415 let first = state.client();
3416 let second = state.client();
3417
3418 assert_eq!(first.app_tag(), expected);
3419 assert!(Arc::ptr_eq(&first, &second));
3420 }
3421
3422 #[test]
3423 fn managed_session_uses_persisted_native_device_id() {
3424 let identity = openrtc::native_device::NativeDeviceIdentity {
3425 device_id: "persisted-native-device".to_string(),
3426 device_name: "Mac".to_string(),
3427 created_at_ms: 1,
3428 updated_at_ms: 1,
3429 name_source: None,
3430 system_info: None,
3431 };
3432
3433 assert_eq!(
3434 managed_session_device_id(Some(" desktop-e2e-native "), &identity),
3435 "persisted-native-device"
3436 );
3437 assert_eq!(
3438 managed_session_device_id(None, &identity),
3439 "persisted-native-device"
3440 );
3441 }
3442
3443 #[test]
3444 fn repeated_managed_session_triggers_converge_on_one_owner() {
3445 let key = "app:persisted-native-device";
3446 let cases = [
3447 ("react-remount", true, true, true, true),
3448 ("hmr", true, true, true, true),
3449 ("auth-refresh", true, true, true, true),
3450 ("resume", true, true, true, true),
3451 ("alias-change", true, true, true, true),
3452 ];
3453 for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
3454 assert_eq!(
3455 managed_session_disposition(
3456 Some(key),
3457 true,
3458 active_presence,
3459 active_auto,
3460 key,
3461 requested_presence,
3462 requested_auto,
3463 ),
3464 SessionDisposition::Reuse,
3465 "{label} must reuse the authoritative tuple"
3466 );
3467 }
3468
3469 assert_eq!(
3470 managed_session_disposition(Some(key), true, false, true, key, true, true),
3471 SessionDisposition::Refresh,
3472 "failed presence startup must retry idempotently"
3473 );
3474 assert_eq!(
3475 managed_session_disposition(
3476 Some(key),
3477 true,
3478 true,
3479 true,
3480 "app:other-native-device",
3481 true,
3482 true,
3483 ),
3484 SessionDisposition::Replace,
3485 "a physical native device change must replace the previous lifecycle owner"
3486 );
3487 assert_eq!(
3488 managed_session_disposition(Some(key), false, true, true, key, true, true),
3489 SessionDisposition::Replace,
3490 "the same tuple on a different client epoch must replace the previous owner"
3491 );
3492 }
3493
3494 #[test]
3495 fn user_device_revocation_invalidates_the_cached_managed_ticket() {
3496 assert!(revokes_managed_session("user-device"));
3497 assert!(revokes_managed_session(" user-device "));
3498 assert!(!revokes_managed_session("share:example"));
3499 assert!(!revokes_managed_session(""));
3500 }
3501
3502 #[test]
3503 fn managed_session_metadata_uses_authoritative_device_id() {
3504 let raw = serde_json::json!({
3505 "deviceId": "persisted-native-device",
3506 "assistantDevice": {
3507 "deviceId": "desktop-e2e-native"
3508 }
3509 })
3510 .to_string();
3511
3512 let metadata =
3513 metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
3514 let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
3515
3516 assert_eq!(parsed["deviceId"], "desktop-e2e-native");
3517 assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
3518 }
3519}