1use std::collections::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 openrtc::native_node::IncomingStreamType;
9use serde::{Deserialize, Serialize};
10use tauri::{Emitter, Manager, Runtime};
11use tokio::sync::Mutex;
12
13const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
14const CONNECTION_STATE_CHANGED_EVENT: &str = "connection-state-changed";
15const PEER_BI_STREAM_INCOMING_EVENT: &str = "openrtc://peer-bi-stream/incoming";
16const PEER_BI_STREAM_CHUNK_EVENT: &str = "openrtc://peer-bi-stream/chunk";
17const PEER_BI_STREAM_CLOSED_EVENT: &str = "openrtc://peer-bi-stream/closed";
18
19pub type InstallFuture = std::pin::Pin<
20 Box<
21 dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
22 + Send,
23 >,
24>;
25
26#[derive(Debug, Clone)]
27pub struct InstallContext {
28 pub data_dir: PathBuf,
32}
33
34pub trait TransportInstaller: Send + Sync {
38 fn id(&self) -> &'static str;
39 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
40 fn install(
41 &self,
42 client: Arc<openrtc::client::Client>,
43 config: openrtc::client::TransportConfig,
44 context: InstallContext,
45 ) -> InstallFuture;
46}
47
48pub trait DeviceKeySigner: Send + Sync {
53 fn public_jwk(&self, app_tag: &str) -> Result<serde_json::Value, String>;
54 fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String>;
55 fn delete(&self, app_tag: &str) -> Result<(), String>;
56 fn read_secure_record(&self, _app_tag: &str, _key: &str) -> Result<Option<String>, String> {
57 Ok(None)
58 }
59 fn write_secure_record(&self, _app_tag: &str, _key: &str, _value: &str) -> Result<(), String> {
60 Err("OpenRTC 2.0 native certificate persistence requires a host secure store".to_string())
61 }
62 fn delete_secure_record(&self, _app_tag: &str, _key: &str) -> Result<(), String> {
63 Ok(())
64 }
65}
66
67struct InstalledNativeTransport {
68 client_ptr: usize,
69 _runtime: Box<dyn std::any::Any + Send + Sync>,
70}
71
72#[derive(Debug, Clone, Serialize, Deserialize)]
73#[serde(rename_all = "camelCase")]
74pub struct OpenRtcTauriConfig {
75 pub api_key: String,
79 #[serde(default)]
80 pub data_dir: Option<PathBuf>,
81 #[serde(default)]
82 pub transport_config: Option<openrtc::client::TransportConfig>,
83}
84
85impl Default for OpenRtcTauriConfig {
86 fn default() -> Self {
87 Self {
88 api_key: String::new(),
89 data_dir: None,
90 transport_config: None,
91 }
92 }
93}
94
95impl OpenRtcTauriConfig {
96 pub fn from_env() -> Self {
97 let api_key = first_env(&[
98 "VITE_OPENRTC_KEY",
99 "VITE_OPENRTC_API_KEY",
100 "VITE_PLUTO_OPENRTC_API_KEY",
101 "OPENRTC_API_KEY",
102 ]);
103 Self {
104 api_key: api_key.unwrap_or_default(),
105 data_dir: None,
106 transport_config: None,
107 }
108 }
109
110 pub fn validated_api_key(&self) -> Result<&str, String> {
111 openrtc::validate_api_key(&self.api_key)
112 .map_err(|error| format!("invalid OpenRTC 2.0 public API key: {error}"))
113 }
114
115 pub fn app_tag(&self) -> Result<String, String> {
116 self.validated_api_key().map(openrtc::app_tag_from_api_key)
117 }
118}
119
120fn first_env(names: &[&str]) -> Option<String> {
121 names.iter().find_map(|name| {
122 std::env::var(name)
123 .ok()
124 .map(|value| value.trim().to_string())
125 .filter(|value| !value.is_empty())
126 })
127}
128
129#[derive(Default)]
130struct TokenRelayState {
131 identity_credential: RwLock<Option<String>>,
132}
133
134impl TokenRelayState {
135 fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
136 let relay = self.clone();
137 Box::new(move || {
138 relay
139 .identity_credential
140 .read()
141 .ok()
142 .and_then(|guard| guard.clone())
143 })
144 }
145
146 fn set(&self, identity_credential: Option<String>) {
147 if let Ok(mut guard) = self.identity_credential.write() {
148 *guard = normalize_token(identity_credential);
149 }
150 }
151}
152
153fn normalize_token(value: Option<String>) -> Option<String> {
154 value
155 .map(|value| value.trim().to_string())
156 .filter(|value| !value.is_empty())
157}
158
159pub struct OpenRtcTauriState {
160 client: RwLock<Arc<openrtc::client::Client>>,
161 config: OpenRtcTauriConfig,
162 token_relay: Arc<TokenRelayState>,
163 data_dir: Option<PathBuf>,
164 connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
165 presence_loop_active: Mutex<bool>,
166 subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
167 peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
168 managed_session_start_guard: Mutex<()>,
169 managed_session_next_owner_epoch: AtomicU64,
170 managed_session: Mutex<Option<ManagedSessionRecord>>,
171 native_transport_installers: Vec<Arc<dyn TransportInstaller>>,
172 installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
173 native_device_key_signer: Option<Arc<dyn DeviceKeySigner>>,
174}
175
176impl OpenRtcTauriState {
177 pub fn new(config: OpenRtcTauriConfig) -> Self {
178 openrtc::ensure_rustls();
179
180 let token_relay = Arc::new(TokenRelayState::default());
181 let client = build_client(&config, &token_relay)
182 .expect("OpenRTC Tauri 2.0 requires a valid public API key");
183 let data_dir = config.data_dir.clone();
184
185 Self {
186 client: RwLock::new(client),
187 config,
188 token_relay,
189 data_dir,
190 connection_state_forwarder: std::sync::Mutex::new(None),
191 presence_loop_active: Mutex::new(false),
192 subscriptions: Mutex::new(HashMap::new()),
193 peer_bi_streams: Mutex::new(HashMap::new()),
194 managed_session_start_guard: Mutex::new(()),
195 managed_session_next_owner_epoch: AtomicU64::new(0),
196 managed_session: Mutex::new(None),
197 native_transport_installers: Vec::new(),
198 installed_native_transports: Mutex::new(HashMap::new()),
199 native_device_key_signer: None,
200 }
201 }
202
203 pub fn with_native_transport_installer(
204 mut self,
205 installer: Arc<dyn TransportInstaller>,
206 ) -> Self {
207 self.native_transport_installers.push(installer);
208 self
209 }
210
211 pub fn with_native_device_key_signer(mut self, signer: Arc<dyn DeviceKeySigner>) -> Self {
212 self.native_device_key_signer = Some(signer);
213 self
214 }
215
216 fn device_public_key(&self, app_tag: &str) -> Result<serde_json::Value, String> {
217 let app_tag = required_app_tag(app_tag)?;
218 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
219 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
220 })?;
221 let value = signer.public_jwk(app_tag)?;
222 validate_public_device_jwk(&value)?;
223 Ok(value)
224 }
225
226 fn sign_device_proof(&self, app_tag: &str, challenge: &str) -> Result<String, String> {
227 let app_tag = required_app_tag(app_tag)?;
228 if challenge.is_empty() || challenge.len() > 2_048 {
229 return Err("OpenRTC device challenge is invalid".to_string());
230 }
231 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
232 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
233 })?;
234 let signature = signer.sign(app_tag, challenge.as_bytes())?;
235 if signature.len() != 64 {
236 return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
237 }
238 Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
239 }
240
241 fn delete_device_key(&self, app_tag: &str) -> Result<(), String> {
242 let app_tag = required_app_tag(app_tag)?;
243 let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
244 "OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
245 })?;
246 signer.delete(app_tag)
247 }
248
249 fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
250 let app_tag = required_app_tag(app_tag)?;
251 let key = required_secure_record_key(key)?;
252 self.native_device_key_signer
253 .as_ref()
254 .ok_or_else(|| {
255 "OpenRTC 2.0 native certificate persistence requires a host secure store"
256 .to_string()
257 })?
258 .read_secure_record(app_tag, key)
259 }
260
261 fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
262 let app_tag = required_app_tag(app_tag)?;
263 let key = required_secure_record_key(key)?;
264 if value.is_empty() || value.len() > 16 * 1024 {
265 return Err("OpenRTC secure record is invalid".to_string());
266 }
267 self.native_device_key_signer
268 .as_ref()
269 .ok_or_else(|| {
270 "OpenRTC 2.0 native certificate persistence requires a host secure store"
271 .to_string()
272 })?
273 .write_secure_record(app_tag, key, value)
274 }
275
276 fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
277 let app_tag = required_app_tag(app_tag)?;
278 let key = required_secure_record_key(key)?;
279 self.native_device_key_signer
280 .as_ref()
281 .ok_or_else(|| {
282 "OpenRTC 2.0 native certificate persistence requires a host secure store"
283 .to_string()
284 })?
285 .delete_secure_record(app_tag, key)
286 }
287
288 pub fn client(&self) -> Arc<openrtc::client::Client> {
289 self.client
290 .read()
291 .map(|guard| guard.clone())
292 .expect("OpenRTC Tauri client state is poisoned")
293 }
294
295 pub fn set_identity_credential(&self, identity_credential: Option<String>) {
296 self.token_relay.set(identity_credential);
297 }
298
299 fn allocate_managed_session_owner_epoch(&self) -> u64 {
300 self.managed_session_next_owner_epoch
301 .fetch_add(1, Ordering::SeqCst)
302 + 1
303 }
304
305 fn replace_connection_state_forwarder<R: Runtime>(
306 &self,
307 app: tauri::AppHandle<R>,
308 client: Arc<openrtc::client::Client>,
309 ) {
310 let next = tauri::async_runtime::spawn(async move {
311 forward_connection_state_events(app, client).await;
312 });
313 if let Ok(mut guard) = self.connection_state_forwarder.lock() {
314 if let Some(previous) = guard.replace(next) {
315 previous.abort();
316 }
317 }
318 }
319}
320
321fn build_client(
322 config: &OpenRtcTauriConfig,
323 token_relay: &Arc<TokenRelayState>,
324) -> Result<Arc<openrtc::client::Client>, String> {
325 let mut builder = openrtc::client::Client::builder(
326 config.validated_api_key()?.to_string(),
327 token_relay.token_provider(),
328 )
329 .map_err(|error| error.to_string())?;
330 if let Some(transport_config) = config.transport_config.clone() {
331 builder = builder.transport_config(transport_config);
332 }
333 Ok(Arc::new(builder.build()))
334}
335
336fn requested_local_device_id(value: Option<&str>) -> Option<String> {
337 value
338 .map(str::trim)
339 .filter(|value| !value.is_empty())
340 .map(ToOwned::to_owned)
341}
342
343fn managed_session_device_id(
344 _requested: Option<&str>,
345 native_identity: &openrtc::native_device::NativeDeviceIdentity,
346) -> String {
347 native_identity.device_id.clone()
351}
352
353#[cfg(test)]
354fn metadata_with_authoritative_device_id(
355 metadata: Option<String>,
356 device_id: &str,
357) -> Option<String> {
358 let device_id = device_id.trim();
359 if device_id.is_empty() {
360 return metadata;
361 }
362
363 let Some(raw_metadata) = metadata else {
364 return Some(serde_json::json!({ "deviceId": device_id }).to_string());
365 };
366
367 match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
368 Ok(serde_json::Value::Object(mut map)) => {
369 map.insert(
370 "deviceId".to_string(),
371 serde_json::Value::String(device_id.to_string()),
372 );
373 Some(serde_json::Value::Object(map).to_string())
374 }
375 _ => Some(
376 serde_json::json!({
377 "deviceId": device_id,
378 "metadata": raw_metadata,
379 })
380 .to_string(),
381 ),
382 }
383}
384
385struct PeerBiStreamHandle {
386 send: Arc<Mutex<Option<PeerSendStream>>>,
387 recv: Option<PeerRecvStream>,
388 read_task: Option<tokio::task::JoinHandle<()>>,
389}
390
391#[derive(Serialize)]
392#[serde(rename_all = "camelCase")]
393pub struct OpenBiResult {
394 stream_id: String,
395 connection_id: Option<String>,
396 remote_node_id: String,
397}
398
399#[derive(Serialize, Clone)]
400#[serde(rename_all = "camelCase")]
401struct IncomingPeerBiStreamEvent {
402 request_id: String,
403 stream_id: String,
404 remote_node_id: String,
405 transport_stable_id: u64,
406}
407
408#[derive(Serialize, Clone)]
409#[serde(rename_all = "camelCase")]
410struct PeerBiStreamChunkEvent {
411 stream_id: String,
412 bytes: Vec<u8>,
413}
414
415#[derive(Serialize, Clone)]
416#[serde(rename_all = "camelCase")]
417struct PeerBiStreamClosedEvent {
418 stream_id: String,
419 error: Option<String>,
420}
421
422#[derive(Debug, Clone, Serialize)]
423#[serde(rename_all = "camelCase")]
424pub struct StartSessionResult {
425 local_node_id: String,
426 ticket_scope: Option<String>,
427 ticket: Option<String>,
428 presence_started: bool,
429 auto_connect_started: bool,
430 local_device: openrtc::native_device::NativeDeviceIdentity,
431}
432
433#[derive(Clone)]
434struct ManagedSessionRecord {
435 key: String,
436 owner_client: Arc<openrtc::client::Client>,
437 owner_epoch: u64,
438 result: StartSessionResult,
439}
440
441#[derive(Debug, Clone, Copy, PartialEq, Eq)]
442enum SessionDisposition {
443 Start,
444 Reuse,
445 Refresh,
446 Replace,
447}
448
449fn managed_session_disposition(
450 active_key: Option<&str>,
451 active_client_matches: bool,
452 active_presence_started: bool,
453 active_auto_connect_started: bool,
454 requested_key: &str,
455 requested_presence: bool,
456 requested_auto_connect: bool,
457) -> SessionDisposition {
458 let Some(active_key) = active_key else {
459 return SessionDisposition::Start;
460 };
461 if !active_client_matches || active_key != requested_key {
462 return SessionDisposition::Replace;
463 }
464 if (!requested_presence || active_presence_started)
465 && (!requested_auto_connect || active_auto_connect_started)
466 {
467 SessionDisposition::Reuse
468 } else {
469 SessionDisposition::Refresh
470 }
471}
472
473fn revokes_managed_session(scope: &str) -> bool {
474 scope.trim() == "user-device"
475}
476
477pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
478 init_with_state(OpenRtcTauriState::new(config))
479}
480
481pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
482 tauri::plugin::Builder::new(PLUGIN_NAME)
483 .setup(move |app, _api| {
484 let client = state.client();
485 let transport_config = state.config.transport_config.clone();
486 let install_context = InstallContext {
491 data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
492 };
493 tauri::async_runtime::block_on(ensure_requested_native_transports(
494 &state,
495 &client,
496 transport_config.as_ref(),
497 &install_context,
498 ))
499 .map_err(std::io::Error::other)?;
500 state.replace_connection_state_forwarder(app.clone(), client);
501 app.manage(state);
502 Ok(())
503 })
504 .invoke_handler(tauri::generate_handler![
505 openrtc_set_identity_credential,
506 openrtc_device_public_key,
507 openrtc_sign_device_proof,
508 openrtc_delete_device_key,
509 openrtc_read_secure_record,
510 openrtc_write_secure_record,
511 openrtc_delete_secure_record,
512 rtc_native_status,
513 get_rtc_local_device_info,
514 update_rtc_local_device_name,
515 get_iroh_node_id,
516 start_iroh_node,
517 get_iroh_endpoint_ticket,
518 register_session_token,
519 get_endpoint_ticket_with_token,
520 validate_session_token,
521 revoke_session_tokens_by_scope,
522 stop_rtc_presence_loop,
523 start_rtc_managed_session,
524 start_rtc_external_auto_connect,
525 submit_rtc_desired_peers,
526 stop_rtc_auto_connect,
527 notify_rtc_network_change,
528 connect_to_device,
529 disconnect_device,
530 set_auto_connect_excluded,
531 set_rtc_external_auto_connect_excluded,
532 resolve_rtc_peer_connection_records,
533 resolve_rtc_peer_identity,
534 get_rtc_peer_session,
535 list_rtc_peer_sessions,
536 list_rtc_managed_connections,
537 wait_for_rtc_settled_peer,
538 list_rtc_connection_states,
539 get_rtc_connection_state,
540 stop_rtc_subscription,
541 start_incoming_peer_bi_streams,
542 is_current_transport_stable_id,
543 open_peer_bi_stream,
544 open_peer_bi_transport_only_stream,
545 open_peer_native_bi_stream,
546 write_peer_bi_stream,
547 start_peer_bi_stream_read,
548 close_peer_bi_stream,
549 ])
550 .build()
551}
552
553pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
554 init(OpenRtcTauriConfig::from_env())
555}
556
557async fn forward_connection_state_events<R: Runtime>(
558 app: tauri::AppHandle<R>,
559 client: Arc<openrtc::client::Client>,
560) {
561 let mut rx = client.connection_state_updates();
562 while let Ok(snapshot) = rx.recv().await {
563 let _ = app.emit(CONNECTION_STATE_CHANGED_EVENT, snapshot);
564 }
565}
566
567fn app_data_dir<R: Runtime>(
568 app: &tauri::AppHandle<R>,
569 state: &OpenRtcTauriState,
570) -> Result<PathBuf, String> {
571 state
572 .data_dir
573 .clone()
574 .or_else(|| app.path().app_data_dir().ok())
575 .ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
576}
577
578async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
579 if let Some(node_id) = client.current_node_id().await {
580 return Ok(node_id);
581 }
582 client
583 .init_iroh(None, Vec::new())
584 .await
585 .map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
586}
587
588async fn ensure_requested_native_transports(
589 state: &OpenRtcTauriState,
590 client: &Arc<openrtc::client::Client>,
591 transports: Option<&openrtc::client::TransportConfig>,
592 context: &InstallContext,
593) -> Result<(), String> {
594 let Some(transports) = transports else {
595 return Ok(());
596 };
597 let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
598 if ble_requested
599 && !state
600 .native_transport_installers
601 .iter()
602 .any(|installer| installer.is_requested(transports))
603 {
604 eprintln!(
605 "[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
606 );
607 return Ok(());
608 }
609 let client_ptr = Arc::as_ptr(client) as usize;
610 for installer in &state.native_transport_installers {
611 if !installer.is_requested(transports) {
612 continue;
613 }
614 let already_installed = state
615 .installed_native_transports
616 .lock()
617 .await
618 .get(installer.id())
619 .is_some_and(|installed| installed.client_ptr == client_ptr);
620 if already_installed {
621 continue;
622 }
623 if client.current_node_id().await.is_some() {
624 eprintln!(
625 "[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
626 installer.id()
627 );
628 continue;
629 }
630 let runtime = match installer
631 .install(client.clone(), transports.clone(), context.clone())
632 .await
633 {
634 Ok(runtime) => runtime,
635 Err(error) => {
636 eprintln!(
637 "[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
638 installer.id(),
639 error
640 );
641 continue;
642 }
643 };
644 state.installed_native_transports.lock().await.insert(
645 installer.id(),
646 InstalledNativeTransport {
647 client_ptr,
648 _runtime: runtime,
649 },
650 );
651 }
652 Ok(())
653}
654
655async fn register_peer_bi_stream<R: Runtime>(
656 app: tauri::AppHandle<R>,
657 state: &OpenRtcTauriState,
658 connection_id: Option<String>,
659 remote_node_id: String,
660 send: PeerSendStream,
661 recv: PeerRecvStream,
662 start_reading: bool,
663) -> OpenBiResult {
664 let stream_id = uuid::Uuid::new_v4().to_string();
665 let send = Arc::new(Mutex::new(Some(send)));
666 let (recv, read_task) = if start_reading {
667 (
668 None,
669 Some(spawn_peer_bi_stream_reader(
670 app.clone(),
671 stream_id.clone(),
672 recv,
673 )),
674 )
675 } else {
676 (Some(recv), None)
680 };
681
682 state.peer_bi_streams.lock().await.insert(
683 stream_id.clone(),
684 PeerBiStreamHandle {
685 send,
686 recv,
687 read_task,
688 },
689 );
690
691 OpenBiResult {
692 stream_id,
693 connection_id,
694 remote_node_id,
695 }
696}
697
698fn spawn_peer_bi_stream_reader<R: Runtime>(
699 app: tauri::AppHandle<R>,
700 stream_id: String,
701 mut recv: PeerRecvStream,
702) -> tokio::task::JoinHandle<()> {
703 tokio::spawn(async move {
704 let mut chunk = vec![0_u8; 64 * 1024];
705 let mut close_error: Option<String> = None;
706 loop {
707 match recv.read(&mut chunk).await {
708 Ok(0) => break,
709 Ok(n) => {
710 if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
711 eprintln!(
712 "[openrtc-tauri][peer-bi-stream] stream_id={} phase=plaintext-chunk bytes={}",
713 stream_id, n
714 );
715 }
716 let _ = app.emit(
717 PEER_BI_STREAM_CHUNK_EVENT,
718 PeerBiStreamChunkEvent {
719 stream_id: stream_id.clone(),
720 bytes: chunk[..n].to_vec(),
721 },
722 );
723 }
724 Err(error) => {
725 if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
726 eprintln!(
727 "[openrtc-tauri][peer-bi-stream] stream_id={} phase=read-error error={}",
728 stream_id, error
729 );
730 }
731 close_error = Some(error.to_string());
732 break;
733 }
734 }
735 }
736 let _ = app.emit(
737 PEER_BI_STREAM_CLOSED_EVENT,
738 PeerBiStreamClosedEvent {
739 stream_id,
740 error: close_error,
741 },
742 );
743 })
744}
745
746async fn open_peer_bi_with<R: Runtime, F, Fut>(
747 app: tauri::AppHandle<R>,
748 state: tauri::State<'_, OpenRtcTauriState>,
749 peer_id: String,
750 timeout_ms: Option<u64>,
751 open: F,
752) -> Result<OpenBiResult, String>
753where
754 F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
755 Fut: std::future::Future<
756 Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
757 >,
758{
759 let peer_id = peer_id.trim().to_string();
760 if peer_id.is_empty() {
761 return Err("peerId is required".to_string());
762 }
763
764 let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
765 .await
766 .map_err(|error| format!("open peer bi stream failed: {error}"))?;
767 Ok(register_peer_bi_stream(app, &state, connection_id, remote_node_id, send, recv, true).await)
768}
769
770#[tauri::command]
774async fn openrtc_set_identity_credential(
775 state: tauri::State<'_, OpenRtcTauriState>,
776 identity_credential: Option<String>,
777) -> Result<(), String> {
778 state.set_identity_credential(identity_credential);
779 if *state.presence_loop_active.lock().await {
780 state.client().request_presence_update();
781 }
782 Ok(())
783}
784
785fn required_app_tag(value: &str) -> Result<&str, String> {
786 let value = value.trim();
787 if value.len() < 5
788 || value.len() > 80
789 || !value.bytes().all(|byte| {
790 byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
791 })
792 {
793 return Err("OpenRTC app tag is invalid".to_string());
794 }
795 Ok(value)
796}
797
798fn validate_public_device_jwk(value: &serde_json::Value) -> Result<(), String> {
799 let object = value
800 .as_object()
801 .ok_or_else(|| "OpenRTC device public key is invalid".to_string())?;
802 let x = object
803 .get("x")
804 .and_then(serde_json::Value::as_str)
805 .unwrap_or_default();
806 if object.get("kty").and_then(serde_json::Value::as_str) != Some("OKP")
807 || object.get("crv").and_then(serde_json::Value::as_str) != Some("Ed25519")
808 || object.contains_key("d")
809 || x.len() != 43
810 || !x
811 .bytes()
812 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
813 {
814 return Err("OpenRTC device public key must be a public Ed25519 JWK".to_string());
815 }
816 Ok(())
817}
818
819fn required_secure_record_key(value: &str) -> Result<&str, String> {
820 let value = value.trim();
821 if value.is_empty()
822 || value.len() > 512
823 || !value
824 .bytes()
825 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':'))
826 {
827 return Err("OpenRTC secure record key is invalid".to_string());
828 }
829 Ok(value)
830}
831
832#[tauri::command]
833async fn openrtc_device_public_key(
834 state: tauri::State<'_, OpenRtcTauriState>,
835 app_tag: String,
836) -> Result<serde_json::Value, String> {
837 state.device_public_key(&app_tag)
838}
839
840#[tauri::command]
841async fn openrtc_sign_device_proof(
842 state: tauri::State<'_, OpenRtcTauriState>,
843 app_tag: String,
844 challenge: String,
845) -> Result<String, String> {
846 state.sign_device_proof(&app_tag, &challenge)
847}
848
849#[tauri::command]
850async fn openrtc_delete_device_key(
851 state: tauri::State<'_, OpenRtcTauriState>,
852 app_tag: String,
853) -> Result<(), String> {
854 state.delete_device_key(&app_tag)
855}
856
857#[tauri::command]
858async fn openrtc_read_secure_record(
859 state: tauri::State<'_, OpenRtcTauriState>,
860 app_tag: String,
861 key: String,
862) -> Result<Option<String>, String> {
863 state.read_secure_record(&app_tag, &key)
864}
865
866#[tauri::command]
867async fn openrtc_write_secure_record(
868 state: tauri::State<'_, OpenRtcTauriState>,
869 app_tag: String,
870 key: String,
871 value: String,
872) -> Result<(), String> {
873 state.write_secure_record(&app_tag, &key, &value)
874}
875
876#[tauri::command]
877async fn openrtc_delete_secure_record(
878 state: tauri::State<'_, OpenRtcTauriState>,
879 app_tag: String,
880 key: String,
881) -> Result<(), String> {
882 state.delete_secure_record(&app_tag, &key)
883}
884
885#[tauri::command]
886async fn rtc_native_status(
887 state: tauri::State<'_, OpenRtcTauriState>,
888) -> Result<openrtc::client::RuntimeStatus, String> {
889 Ok(state.client().runtime_status().await)
890}
891
892#[tauri::command]
893async fn get_rtc_local_device_info<R: Runtime>(
894 app: tauri::AppHandle<R>,
895 state: tauri::State<'_, OpenRtcTauriState>,
896) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
897 let data_dir = app_data_dir(&app, &state)?;
898 state
899 .client()
900 .init_native_device_identity(data_dir, None)
901 .await
902 .map_err(|error| error.to_string())
903}
904
905#[tauri::command]
906async fn update_rtc_local_device_name<R: Runtime>(
907 app: tauri::AppHandle<R>,
908 state: tauri::State<'_, OpenRtcTauriState>,
909 device_name: String,
910) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
911 let data_dir = app_data_dir(&app, &state)?;
912 let _ = state
913 .client()
914 .init_native_device_identity(data_dir, None)
915 .await
916 .map_err(|error| error.to_string())?;
917 state
918 .client()
919 .update_native_device_name(&device_name)
920 .await
921 .map_err(|error| error.to_string())
922}
923
924#[tauri::command]
925async fn get_iroh_node_id(
926 state: tauri::State<'_, OpenRtcTauriState>,
927) -> Result<Option<String>, String> {
928 Ok(state.client().current_node_id().await)
929}
930
931#[tauri::command]
932async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
933 ensure_iroh_node(state.client()).await
934}
935
936#[tauri::command]
937async fn get_iroh_endpoint_ticket(
938 state: tauri::State<'_, OpenRtcTauriState>,
939) -> Result<String, String> {
940 ensure_iroh_node(state.client()).await?;
941 state
942 .client()
943 .endpoint_ticket()
944 .await
945 .map_err(|error| error.to_string())
946}
947
948#[tauri::command]
949async fn register_session_token(
950 state: tauri::State<'_, OpenRtcTauriState>,
951 token: String,
952 scope: String,
953 max_connections: u32,
954 expires_at_ms: Option<u64>,
955) -> Result<(), String> {
956 if let Some(expires_at_ms) = expires_at_ms {
957 state
958 .client()
959 .register_token_until(token, scope, max_connections, expires_at_ms);
960 } else {
961 state
962 .client()
963 .register_session_token(token, scope, max_connections);
964 }
965 Ok(())
966}
967
968#[tauri::command]
969async fn get_endpoint_ticket_with_token(
970 state: tauri::State<'_, OpenRtcTauriState>,
971 scope: String,
972 max_connections: u32,
973) -> Result<String, String> {
974 ensure_iroh_node(state.client()).await?;
975 state
976 .client()
977 .endpoint_ticket_with_token(&scope, max_connections)
978 .await
979 .map_err(|error| error.to_string())
980}
981
982#[tauri::command]
983async fn validate_session_token(
984 state: tauri::State<'_, OpenRtcTauriState>,
985 token: String,
986 connection_id: Option<String>,
987) -> Result<String, String> {
988 if let Some(connection_id) = connection_id
989 .as_deref()
990 .map(str::trim)
991 .filter(|value| !value.is_empty())
992 {
993 state
994 .client()
995 .validate_connection_token(&token, connection_id, None)
996 .await
997 .map_err(|error| error.to_string())
998 } else {
999 state.client().validate_session_token(&token)
1000 }
1001}
1002
1003#[tauri::command]
1004async fn revoke_session_tokens_by_scope(
1005 state: tauri::State<'_, OpenRtcTauriState>,
1006 scope: String,
1007) -> Result<Vec<String>, String> {
1008 let affected = state.client().revoke_tokens_by_scope(&scope).await;
1009 if revokes_managed_session(&scope) {
1010 if let Some(previous) = state.managed_session.lock().await.take() {
1017 previous.owner_client.stop_presence_loop();
1018 previous.owner_client.stop_auto_connect();
1019 previous.owner_client.stop_external_auto_connect().await;
1020 }
1021 *state.presence_loop_active.lock().await = false;
1022 }
1023 Ok(affected)
1024}
1025
1026#[tauri::command]
1027async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1028 state.client().stop_presence_loop();
1029 *state.presence_loop_active.lock().await = false;
1030 if let Some(active) = state.managed_session.lock().await.as_mut() {
1031 active.result.presence_started = false;
1032 }
1033 stop_subscription_by_id(&state, "internal-presence-loop").await;
1034 Ok(())
1035}
1036
1037#[tauri::command]
1038async fn start_rtc_managed_session<R: Runtime>(
1039 app: tauri::AppHandle<R>,
1040 state: tauri::State<'_, OpenRtcTauriState>,
1041 user_id: String,
1042 device_name: Option<String>,
1043 local_device_id: Option<String>,
1044 metadata: Option<String>,
1045 transports: Option<openrtc::client::TransportConfig>,
1046 auto_connect: Option<bool>,
1047 presence: Option<bool>,
1048) -> Result<StartSessionResult, String> {
1049 let command_started = std::time::Instant::now();
1050 let _start_guard = state.managed_session_start_guard.lock().await;
1051 let client = state.client();
1052 let app_tag = client.app_tag().to_string();
1053 eprintln!(
1054 "[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
1055 user_id,
1056 app_tag,
1057 device_name
1058 .as_deref()
1059 .map(str::trim)
1060 .map(|value| !value.is_empty())
1061 .unwrap_or(false),
1062 local_device_id
1063 .as_deref()
1064 .map(str::trim)
1065 .map(|value| !value.is_empty())
1066 .unwrap_or(false),
1067 metadata
1068 .as_deref()
1069 .map(str::trim)
1070 .map(|value| !value.is_empty())
1071 .unwrap_or(false),
1072 auto_connect.unwrap_or(false),
1073 presence.unwrap_or(false)
1074 );
1075 let data_dir = app_data_dir(&app, &state)?;
1076 let install_context = InstallContext {
1077 data_dir: data_dir.clone(),
1078 };
1079 ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
1080 .await?;
1081 if let Some(transport_config) = transports {
1082 client
1083 .update_transport_config(transport_config)
1084 .await
1085 .map_err(|error| error.to_string())?;
1086 }
1087
1088 let local_node_id = ensure_iroh_node(client.clone()).await?;
1089 let local_device = client
1090 .init_native_device_identity(data_dir, device_name.as_deref())
1091 .await
1092 .map_err(|error| error.to_string())?;
1093
1094 let effective_local_device_id =
1095 managed_session_device_id(local_device_id.as_deref(), &local_device);
1096 if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
1097 if requested != local_device.device_id {
1098 eprintln!(
1099 "[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
1100 requested, local_device.device_id
1101 );
1102 }
1103 }
1104
1105 let should_start_presence = presence.unwrap_or(false);
1106 let should_start_auto_connect = auto_connect.unwrap_or(false);
1107 if should_start_presence || should_start_auto_connect {
1108 return Err(
1109 "native Firebase coordination has been removed; publish gateway presence and submit external desired peers instead"
1110 .to_string(),
1111 );
1112 }
1113 let managed_session_key = format!("{}:{}:{}", app_tag, user_id, effective_local_device_id);
1114 let (disposition, reused) = {
1115 let active = state.managed_session.lock().await;
1116 let disposition = managed_session_disposition(
1117 active.as_ref().map(|record| record.key.as_str()),
1118 active
1119 .as_ref()
1120 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1121 active
1122 .as_ref()
1123 .is_some_and(|record| record.result.presence_started),
1124 active
1125 .as_ref()
1126 .is_some_and(|record| record.result.auto_connect_started),
1127 &managed_session_key,
1128 should_start_presence,
1129 should_start_auto_connect,
1130 );
1131 let reused = (disposition == SessionDisposition::Reuse).then(|| {
1132 active
1133 .as_ref()
1134 .expect("managed session reuse requires an active record")
1135 .result
1136 .clone()
1137 });
1138 (disposition, reused)
1139 };
1140 if let Some(reused) = reused {
1141 eprintln!(
1142 "[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
1143 managed_session_key,
1144 command_started.elapsed().as_millis()
1145 );
1146 return Ok(reused);
1147 }
1148
1149 if disposition == SessionDisposition::Replace {
1152 let previous = state
1153 .managed_session
1154 .lock()
1155 .await
1156 .take()
1157 .expect("managed session replacement requires an active record");
1158 previous.owner_client.stop_presence_loop();
1159 previous.owner_client.stop_auto_connect();
1160 previous.owner_client.stop_external_auto_connect().await;
1161 *state.presence_loop_active.lock().await = false;
1162 }
1163 let ticket_scope = Some("user-device".to_string());
1164 let ticket = client
1165 .endpoint_ticket_with_token("user-device", 0)
1166 .await
1167 .map_err(|error| error.to_string())?;
1168
1169 let logical_local_device = local_device.clone();
1170
1171 let owner_epoch = match disposition {
1172 SessionDisposition::Refresh => {
1173 state
1174 .managed_session
1175 .lock()
1176 .await
1177 .as_ref()
1178 .expect("managed session refresh requires an active record")
1179 .owner_epoch
1180 }
1181 SessionDisposition::Start | SessionDisposition::Replace => {
1182 state.allocate_managed_session_owner_epoch()
1183 }
1184 SessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
1185 };
1186
1187 eprintln!(
1188 "[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={}",
1189 user_id,
1190 app_tag,
1191 local_node_id,
1192 logical_local_device.device_id,
1193 owner_epoch,
1194 should_start_presence,
1195 should_start_auto_connect,
1196 ticket_scope.as_deref().unwrap_or("unrestricted"),
1197 command_started.elapsed().as_millis()
1198 );
1199
1200 let result = StartSessionResult {
1201 local_node_id,
1202 ticket_scope,
1203 ticket: Some(ticket),
1204 presence_started: false,
1205 auto_connect_started: false,
1206 local_device: logical_local_device,
1207 };
1208 *state.managed_session.lock().await = Some(ManagedSessionRecord {
1209 key: managed_session_key,
1210 owner_client: client,
1211 owner_epoch,
1212 result: result.clone(),
1213 });
1214 Ok(result)
1215}
1216
1217#[tauri::command]
1218async fn start_rtc_external_auto_connect(
1219 state: tauri::State<'_, OpenRtcTauriState>,
1220 user_id: String,
1221 local_device_id: String,
1222) -> Result<(), String> {
1223 state
1224 .client()
1225 .start_external_auto_connect(user_id, local_device_id)
1226 .await
1227 .map_err(|error| error.to_string())
1228}
1229
1230#[tauri::command]
1231async fn submit_rtc_desired_peers(
1232 state: tauri::State<'_, OpenRtcTauriState>,
1233 revision: u64,
1234 peers_json: String,
1235) -> Result<bool, String> {
1236 state
1237 .client()
1238 .submit_external_desired_peers(revision, &peers_json)
1239 .await
1240 .map_err(|error| error.to_string())
1241}
1242
1243#[tauri::command]
1244async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
1245 state.client().stop_auto_connect();
1246 state.client().stop_external_auto_connect().await;
1247 if let Some(active) = state.managed_session.lock().await.as_mut() {
1248 active.result.auto_connect_started = false;
1249 }
1250 Ok(())
1251}
1252
1253#[tauri::command]
1254async fn connect_to_device(
1255 state: tauri::State<'_, OpenRtcTauriState>,
1256 device_id: Option<String>,
1257 endpoint_ticket: String,
1258 timeout_ms: Option<u64>,
1259) -> Result<openrtc::client::ManagedConnectResult, String> {
1260 let client = state.client();
1261 let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
1262 match timeout_ms {
1263 Some(timeout_ms) => {
1264 tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
1265 .await
1266 .map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
1267 .map_err(|error| error.to_string())
1268 }
1269 None => connect.await.map_err(|error| error.to_string()),
1270 }
1271}
1272
1273#[tauri::command]
1274async fn disconnect_device(
1275 state: tauri::State<'_, OpenRtcTauriState>,
1276 device_id: String,
1277 node_id_hint: Option<String>,
1278) -> Result<(), String> {
1279 state
1280 .client()
1281 .disconnect_device(&device_id, node_id_hint.as_deref())
1282 .await;
1283 Ok(())
1284}
1285
1286#[tauri::command]
1287async fn set_auto_connect_excluded(
1288 state: tauri::State<'_, OpenRtcTauriState>,
1289 device_id: String,
1290 excluded: bool,
1291) -> Result<(), String> {
1292 if excluded {
1293 state.client().exclude_peer_and_publish(&device_id).await;
1294 } else {
1295 state.client().unexclude_peer_and_publish(&device_id).await;
1296 }
1297 Ok(())
1298}
1299
1300#[tauri::command]
1301async fn set_rtc_external_auto_connect_excluded(
1302 state: tauri::State<'_, OpenRtcTauriState>,
1303 device_id: String,
1304 excluded: bool,
1305) -> Result<(), String> {
1306 state
1307 .client()
1308 .set_auto_connect_excluded(&device_id, excluded);
1309 state.client().wake_native_external_auto_connect().await;
1310 Ok(())
1311}
1312
1313#[tauri::command]
1314async fn resolve_rtc_peer_connection_records(
1315 state: tauri::State<'_, OpenRtcTauriState>,
1316 id: String,
1317) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1318 Ok(state.client().resolve_peer_connection_records(&id).await)
1319}
1320
1321#[tauri::command]
1322async fn resolve_rtc_peer_identity(
1323 state: tauri::State<'_, OpenRtcTauriState>,
1324 id: String,
1325) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
1326 Ok(state.client().peer_snapshot(&id).await)
1327}
1328
1329#[tauri::command]
1330async fn get_rtc_peer_session(
1331 state: tauri::State<'_, OpenRtcTauriState>,
1332 id: String,
1333) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1334 Ok(state.client().peer_session(&id).await)
1335}
1336
1337#[tauri::command]
1338async fn list_rtc_peer_sessions(
1339 state: tauri::State<'_, OpenRtcTauriState>,
1340) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
1341 Ok(state.client().peer_sessions().await)
1342}
1343
1344#[tauri::command]
1345async fn list_rtc_managed_connections(
1346 state: tauri::State<'_, OpenRtcTauriState>,
1347) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
1348 Ok(state.client().list_managed_connections().await)
1349}
1350
1351#[derive(Debug, Clone, Serialize)]
1352#[serde(rename_all = "camelCase")]
1353struct NotifyRtcNetworkChangeResult {
1354 retired_stale_connections: usize,
1355}
1356
1357#[tauri::command]
1358async fn notify_rtc_network_change(
1359 state: tauri::State<'_, OpenRtcTauriState>,
1360) -> Result<NotifyRtcNetworkChangeResult, String> {
1361 let retired_stale_connections = state
1362 .client()
1363 .notify_network_change()
1364 .await
1365 .map_err(|error| error.to_string())?;
1366 Ok(NotifyRtcNetworkChangeResult {
1367 retired_stale_connections,
1368 })
1369}
1370
1371#[tauri::command]
1372async fn wait_for_rtc_settled_peer(
1373 state: tauri::State<'_, OpenRtcTauriState>,
1374 id: String,
1375 timeout_ms: Option<u64>,
1376) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
1377 Ok(state.client().wait_for_peer(&id, timeout_ms).await)
1378}
1379
1380#[tauri::command]
1381async fn list_rtc_connection_states(
1382 state: tauri::State<'_, OpenRtcTauriState>,
1383) -> Result<Vec<openrtc::client::StateSnapshot>, String> {
1384 Ok(state.client().connection_states().await)
1385}
1386
1387#[tauri::command]
1388async fn get_rtc_connection_state(
1389 state: tauri::State<'_, OpenRtcTauriState>,
1390 connection_id: String,
1391) -> Result<Option<openrtc::client::StateSnapshot>, String> {
1392 Ok(state.client().connection_state(&connection_id).await)
1393}
1394
1395#[tauri::command]
1396async fn stop_rtc_subscription(
1397 state: tauri::State<'_, OpenRtcTauriState>,
1398 request_id: String,
1399) -> Result<(), String> {
1400 stop_subscription_by_id(&state, &request_id).await;
1401 Ok(())
1402}
1403
1404async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
1405 if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
1406 handle.abort();
1407 }
1408}
1409
1410#[tauri::command]
1411async fn start_incoming_peer_bi_streams<R: Runtime>(
1412 app: tauri::AppHandle<R>,
1413 state: tauri::State<'_, OpenRtcTauriState>,
1414 request_id: Option<String>,
1415) -> Result<String, String> {
1416 let request_id = request_id
1417 .map(|value| value.trim().to_string())
1418 .filter(|value| !value.is_empty())
1419 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1420
1421 stop_subscription_by_id(&state, &request_id).await;
1422
1423 let incoming = state
1424 .client()
1425 .incoming_streams()
1426 .await
1427 .map_err(|error| format!("failed to subscribe to incoming peer streams: {error}"))?;
1428 let app_for_task = app.clone();
1429 let request_for_task = request_id.clone();
1430 let handle = tokio::spawn(async move {
1431 while let Ok(incoming_stream) = incoming.recv().await {
1432 let endpoint_id = incoming_stream.endpoint_id;
1433 let remote_node_id = endpoint_id.to_string();
1434 let transport_stable_id = incoming_stream.transport_stable_id;
1435 let recv_prefix = incoming_stream.recv_prefix;
1436 let IncomingStreamType::Bi(send, recv) = incoming_stream.stream else {
1437 continue;
1438 };
1439 if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1440 eprintln!(
1441 "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} request_id={} phase=received",
1442 remote_node_id,
1443 transport_stable_id,
1444 request_for_task
1445 );
1446 }
1447 let state_for_task = app_for_task.state::<OpenRtcTauriState>();
1448 let (send, recv) = match state_for_task
1449 .client()
1450 .wrap_bi_with_prefix(&endpoint_id, transport_stable_id, send, recv, &recv_prefix)
1451 .await
1452 {
1453 Ok(streams) => streams,
1454 Err(error) => {
1455 eprintln!(
1456 "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} phase=rejected error={}",
1457 remote_node_id, transport_stable_id, error
1458 );
1459 continue;
1460 }
1461 };
1462 let result = register_peer_bi_stream(
1463 app_for_task.clone(),
1464 state_for_task.inner(),
1465 None,
1466 remote_node_id.clone(),
1467 send,
1468 recv,
1469 false,
1470 )
1471 .await;
1472 let _ = app_for_task.emit(
1473 PEER_BI_STREAM_INCOMING_EVENT,
1474 IncomingPeerBiStreamEvent {
1475 request_id: request_for_task.clone(),
1476 stream_id: result.stream_id,
1477 remote_node_id: remote_node_id.clone(),
1478 transport_stable_id,
1479 },
1480 );
1481 if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1482 eprintln!(
1483 "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} request_id={} phase=emitted",
1484 remote_node_id,
1485 transport_stable_id,
1486 request_for_task
1487 );
1488 }
1489 }
1490 });
1491
1492 state
1493 .subscriptions
1494 .lock()
1495 .await
1496 .insert(request_id.clone(), handle);
1497 Ok(request_id)
1498}
1499
1500#[tauri::command]
1501async fn is_current_transport_stable_id(
1502 state: tauri::State<'_, OpenRtcTauriState>,
1503 endpoint_id: String,
1504 transport_stable_id: u64,
1505) -> Result<bool, String> {
1506 let is_current = state
1507 .client()
1508 .is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
1509 .await
1510 .map_err(|error| error.to_string())?;
1511 if std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1") {
1512 eprintln!(
1513 "[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
1514 endpoint_id, transport_stable_id, is_current
1515 );
1516 }
1517 Ok(is_current)
1518}
1519
1520#[tauri::command]
1521async fn open_peer_bi_stream<R: Runtime>(
1522 app: tauri::AppHandle<R>,
1523 state: tauri::State<'_, OpenRtcTauriState>,
1524 peer_id: String,
1525 timeout_ms: Option<u64>,
1526) -> Result<OpenBiResult, String> {
1527 open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
1528 client.open_peer_bi(&peer_id, timeout_ms).await
1529 })
1530 .await
1531}
1532
1533#[tauri::command]
1534async fn open_peer_bi_transport_only_stream<R: Runtime>(
1535 app: tauri::AppHandle<R>,
1536 state: tauri::State<'_, OpenRtcTauriState>,
1537 peer_id: String,
1538 timeout_ms: Option<u64>,
1539) -> Result<OpenBiResult, String> {
1540 open_peer_bi_with(
1541 app,
1542 state,
1543 peer_id,
1544 timeout_ms,
1545 |client, peer_id, timeout_ms| async move {
1546 client
1547 .open_peer_bi_transport_only(&peer_id, timeout_ms)
1548 .await
1549 .map(|(connection_id, remote_node_id, send, recv)| {
1550 (
1551 connection_id,
1552 remote_node_id,
1553 PeerSendStream::plain(send),
1554 PeerRecvStream::plain(recv),
1555 )
1556 })
1557 },
1558 )
1559 .await
1560}
1561
1562#[tauri::command]
1563async fn open_peer_native_bi_stream<R: Runtime>(
1564 app: tauri::AppHandle<R>,
1565 state: tauri::State<'_, OpenRtcTauriState>,
1566 peer_id: String,
1567 label: String,
1568 timeout_ms: Option<u64>,
1569) -> Result<OpenBiResult, String> {
1570 let channel_envelope = openrtc::stream_metadata::encode_envelope(&label, None)
1571 .map_err(|error| error.to_string())?;
1572 open_peer_bi_with(
1573 app,
1574 state,
1575 peer_id,
1576 timeout_ms,
1577 move |client, peer_id, timeout_ms| async move {
1578 let (connection_id, remote_node_id, mut send, recv) =
1579 client.open_peer_bi(&peer_id, timeout_ms).await?;
1580 send.write_all(&channel_envelope).await?;
1581 Ok((connection_id, remote_node_id, send, recv))
1582 },
1583 )
1584 .await
1585}
1586
1587#[tauri::command]
1588async fn write_peer_bi_stream(
1589 state: tauri::State<'_, OpenRtcTauriState>,
1590 stream_id: String,
1591 bytes: Vec<u8>,
1592) -> Result<(), String> {
1593 let stream_id = stream_id.trim().to_string();
1594 if stream_id.is_empty() {
1595 return Err("streamId is required".to_string());
1596 }
1597
1598 let send = {
1599 let streams = state.peer_bi_streams.lock().await;
1600 streams
1601 .get(&stream_id)
1602 .map(|handle| handle.send.clone())
1603 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
1604 };
1605 let result = send
1606 .lock()
1607 .await
1608 .as_mut()
1609 .ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
1610 .write_all(&bytes)
1611 .await
1612 .map_err(|error| format!("write peer bi stream failed: {error}"));
1613 result
1614}
1615
1616#[tauri::command]
1617async fn start_peer_bi_stream_read<R: Runtime>(
1618 app: tauri::AppHandle<R>,
1619 state: tauri::State<'_, OpenRtcTauriState>,
1620 stream_id: String,
1621) -> Result<(), String> {
1622 let stream_id = stream_id.trim().to_string();
1623 if stream_id.is_empty() {
1624 return Err("streamId is required".to_string());
1625 }
1626
1627 let mut streams = state.peer_bi_streams.lock().await;
1628 let handle = streams
1629 .get_mut(&stream_id)
1630 .ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
1631 if handle.read_task.is_some() {
1632 return Ok(());
1633 }
1634 let recv = handle
1635 .recv
1636 .take()
1637 .ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
1638 handle.read_task = Some(spawn_peer_bi_stream_reader(app, stream_id, recv));
1639 Ok(())
1640}
1641
1642#[tauri::command]
1643async fn close_peer_bi_stream(
1644 state: tauri::State<'_, OpenRtcTauriState>,
1645 stream_id: String,
1646) -> Result<(), String> {
1647 let stream_id = stream_id.trim().to_string();
1648 if stream_id.is_empty() {
1649 return Err("streamId is required".to_string());
1650 }
1651
1652 let handle = {
1653 let mut streams = state.peer_bi_streams.lock().await;
1654 streams.remove(&stream_id)
1655 };
1656 if let Some(handle) = handle {
1657 let mut send = handle.send.lock().await;
1660 if let Some(send) = send.take() {
1661 let _ = send.finish();
1662 }
1663 drop(send);
1664 if let Some(read_task) = handle.read_task {
1665 read_task.abort();
1666 }
1667 }
1668 Ok(())
1669}
1670
1671#[cfg(test)]
1672mod tests {
1673 use super::*;
1674 use std::collections::HashMap as StdHashMap;
1675 use std::sync::atomic::{AtomicUsize, Ordering};
1676 use std::sync::Mutex as StdMutex;
1677 use tokio::sync::oneshot;
1678
1679 struct FakeTransportInstaller {
1680 installs: Arc<AtomicUsize>,
1681 }
1682
1683 #[derive(Default)]
1684 struct FakeDeviceKeySigner {
1685 records: StdMutex<StdHashMap<String, String>>,
1686 }
1687
1688 impl DeviceKeySigner for FakeDeviceKeySigner {
1689 fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
1690 Ok(serde_json::json!({
1691 "kty": "OKP",
1692 "crv": "Ed25519",
1693 "x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
1694 }))
1695 }
1696
1697 fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
1698 assert_eq!(challenge, b"openrtc:v2:test");
1699 Ok(vec![7; 64])
1700 }
1701
1702 fn delete(&self, _app_tag: &str) -> Result<(), String> {
1703 Ok(())
1704 }
1705
1706 fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
1707 Ok(self
1708 .records
1709 .lock()
1710 .unwrap()
1711 .get(&format!("{app_tag}:{key}"))
1712 .cloned())
1713 }
1714
1715 fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
1716 self.records
1717 .lock()
1718 .unwrap()
1719 .insert(format!("{app_tag}:{key}"), value.to_string());
1720 Ok(())
1721 }
1722
1723 fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
1724 self.records
1725 .lock()
1726 .unwrap()
1727 .remove(&format!("{app_tag}:{key}"));
1728 Ok(())
1729 }
1730 }
1731
1732 impl TransportInstaller for FakeTransportInstaller {
1733 fn id(&self) -> &'static str {
1734 "ble"
1735 }
1736
1737 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
1738 config.ble.as_ref().is_some_and(|ble| ble.enabled)
1739 }
1740
1741 fn install(
1742 &self,
1743 _client: Arc<openrtc::client::Client>,
1744 _config: openrtc::client::TransportConfig,
1745 _context: InstallContext,
1746 ) -> InstallFuture {
1747 self.installs.fetch_add(1, Ordering::SeqCst);
1748 Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
1749 }
1750 }
1751
1752 #[derive(Debug)]
1753 struct UnavailableTransportInstaller;
1754
1755 impl TransportInstaller for UnavailableTransportInstaller {
1756 fn id(&self) -> &'static str {
1757 "ble"
1758 }
1759
1760 fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
1761 config.ble.as_ref().is_some_and(|ble| ble.enabled)
1762 }
1763
1764 fn install(
1765 &self,
1766 _client: Arc<openrtc::client::Client>,
1767 _config: openrtc::client::TransportConfig,
1768 _context: InstallContext,
1769 ) -> InstallFuture {
1770 Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
1771 }
1772 }
1773
1774 fn requested_ble_config() -> openrtc::client::TransportConfig {
1775 openrtc::client::TransportConfig {
1776 ble: Some(openrtc::client::BleConfig {
1777 enabled: true,
1778 ..openrtc::client::BleConfig::default()
1779 }),
1780 ..openrtc::client::TransportConfig::default()
1781 }
1782 }
1783
1784 fn test_config() -> OpenRtcTauriConfig {
1785 OpenRtcTauriConfig {
1786 api_key: format!("pk_test_{}", "a".repeat(40)),
1787 ..OpenRtcTauriConfig::default()
1788 }
1789 }
1790
1791 #[test]
1792 fn native_device_proof_uses_host_signer_without_exporting_private_key() {
1793 let state = OpenRtcTauriState::new(test_config())
1794 .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
1795 let public = state
1796 .device_public_key("app_native_test")
1797 .expect("public key");
1798 assert_eq!(public["kty"], "OKP");
1799 assert_eq!(public["crv"], "Ed25519");
1800 assert!(public.get("d").is_none());
1801 let signature = state
1802 .sign_device_proof("app_native_test", "openrtc:v2:test")
1803 .expect("signature");
1804 assert_eq!(
1805 base64::engine::general_purpose::URL_SAFE_NO_PAD
1806 .decode(signature)
1807 .expect("base64"),
1808 vec![7; 64]
1809 );
1810 }
1811
1812 #[test]
1813 fn native_certificate_records_round_trip_through_host_secure_storage() {
1814 let state = OpenRtcTauriState::new(test_config())
1815 .with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
1816 let app_tag = "app_native_test";
1817 let key = "openrtc:v2:device-session:abc:device-1";
1818 assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
1819 state
1820 .write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
1821 .unwrap();
1822 assert_eq!(
1823 state.read_secure_record(app_tag, key).unwrap().as_deref(),
1824 Some("{\"token\":\"bound\"}")
1825 );
1826 state.delete_secure_record(app_tag, key).unwrap();
1827 assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
1828 }
1829
1830 #[test]
1831 fn native_device_proof_fails_closed_without_secure_host_signer() {
1832 let state = OpenRtcTauriState::new(test_config());
1833 assert!(state
1834 .device_public_key("app_native_test")
1835 .expect_err("missing signer must fail")
1836 .contains("secure-storage signer"));
1837 }
1838
1839 fn native_transport_install_context() -> InstallContext {
1840 InstallContext {
1841 data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
1842 }
1843 }
1844
1845 fn managed_session_test_result(user_id: &str) -> StartSessionResult {
1846 StartSessionResult {
1847 local_node_id: format!("node-{user_id}"),
1848 ticket_scope: Some("user-device".to_string()),
1849 ticket: None,
1850 presence_started: true,
1851 auto_connect_started: true,
1852 local_device: openrtc::native_device::NativeDeviceIdentity {
1853 device_id: format!("device-{user_id}"),
1854 device_name: "Test Device".to_string(),
1855 created_at_ms: 1,
1856 updated_at_ms: 1,
1857 name_source: None,
1858 system_info: None,
1859 },
1860 }
1861 }
1862
1863 async fn simulate_managed_session_start(
1864 state: Arc<OpenRtcTauriState>,
1865 user_id: String,
1866 pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
1867 starts: Arc<AtomicUsize>,
1868 stopped_owners: Arc<Mutex<Vec<usize>>>,
1869 ) -> SessionDisposition {
1870 let _start_guard = state.managed_session_start_guard.lock().await;
1871 let client = state.client();
1872 let key = format!("{}:{user_id}:device-{user_id}", client.app_tag());
1873 let (disposition, active_epoch) = {
1874 let active = state.managed_session.lock().await;
1875 let disposition = managed_session_disposition(
1876 active.as_ref().map(|record| record.key.as_str()),
1877 active
1878 .as_ref()
1879 .is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
1880 active
1881 .as_ref()
1882 .is_some_and(|record| record.result.presence_started),
1883 active
1884 .as_ref()
1885 .is_some_and(|record| record.result.auto_connect_started),
1886 &key,
1887 true,
1888 true,
1889 );
1890 (
1891 disposition,
1892 active.as_ref().map(|record| record.owner_epoch),
1893 )
1894 };
1895
1896 if disposition == SessionDisposition::Reuse {
1897 return disposition;
1898 }
1899
1900 if disposition == SessionDisposition::Replace {
1901 let previous = state
1902 .managed_session
1903 .lock()
1904 .await
1905 .take()
1906 .expect("replacement must have a previous owner");
1907 stopped_owners
1908 .lock()
1909 .await
1910 .push(Arc::as_ptr(&previous.owner_client) as usize);
1911 }
1912
1913 starts.fetch_add(1, Ordering::SeqCst);
1914 if let Some((entered, release)) = pause {
1915 entered.send(()).expect("start observer must be waiting");
1916 release.await.expect("start release must be sent");
1917 }
1918
1919 let owner_epoch = match disposition {
1920 SessionDisposition::Refresh => {
1921 active_epoch.expect("refresh must preserve the active owner epoch")
1922 }
1923 SessionDisposition::Start | SessionDisposition::Replace => {
1924 state.allocate_managed_session_owner_epoch()
1925 }
1926 SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
1927 };
1928 let result = managed_session_test_result(&user_id);
1929 *state.managed_session.lock().await = Some(ManagedSessionRecord {
1930 key,
1931 owner_client: client,
1932 owner_epoch,
1933 result,
1934 });
1935 disposition
1936 }
1937
1938 #[tokio::test]
1939 async fn concurrent_same_key_start_reuses_one_owner_epoch() {
1940 let state = Arc::new(OpenRtcTauriState::new(test_config()));
1941 let starts = Arc::new(AtomicUsize::new(0));
1942 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
1943 let (first_entered_tx, first_entered_rx) = oneshot::channel();
1944 let (first_release_tx, first_release_rx) = oneshot::channel();
1945
1946 let first = tokio::spawn(simulate_managed_session_start(
1947 state.clone(),
1948 "same-user".to_string(),
1949 Some((first_entered_tx, first_release_rx)),
1950 starts.clone(),
1951 stopped_owners.clone(),
1952 ));
1953 first_entered_rx.await.expect("first start must pause");
1954
1955 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
1956 let second_state = state.clone();
1957 let second_starts = starts.clone();
1958 let second_stopped_owners = stopped_owners.clone();
1959 let second = tokio::spawn(async move {
1960 second_attempting_tx.send(()).unwrap();
1961 simulate_managed_session_start(
1962 second_state,
1963 "same-user".to_string(),
1964 None,
1965 second_starts,
1966 second_stopped_owners,
1967 )
1968 .await
1969 });
1970 second_attempting_rx.await.unwrap();
1971 tokio::task::yield_now().await;
1972 assert_eq!(starts.load(Ordering::SeqCst), 1);
1973
1974 first_release_tx.send(()).unwrap();
1975 assert_eq!(
1976 first.await.expect("first start task"),
1977 SessionDisposition::Start
1978 );
1979 assert_eq!(
1980 second.await.expect("second start task"),
1981 SessionDisposition::Reuse
1982 );
1983 assert_eq!(starts.load(Ordering::SeqCst), 1);
1984
1985 let active = state
1986 .managed_session
1987 .lock()
1988 .await
1989 .clone()
1990 .expect("active owner");
1991 assert_eq!(active.owner_epoch, 1);
1992 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
1993 assert!(stopped_owners.lock().await.is_empty());
1994 }
1995
1996 #[tokio::test]
1997 async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
1998 let state = Arc::new(OpenRtcTauriState::new(test_config()));
1999 let starts = Arc::new(AtomicUsize::new(0));
2000 let stopped_owners = Arc::new(Mutex::new(Vec::new()));
2001 let (first_entered_tx, first_entered_rx) = oneshot::channel();
2002 let (first_release_tx, first_release_rx) = oneshot::channel();
2003
2004 let first = tokio::spawn(simulate_managed_session_start(
2005 state.clone(),
2006 "first-user".to_string(),
2007 Some((first_entered_tx, first_release_rx)),
2008 starts.clone(),
2009 stopped_owners.clone(),
2010 ));
2011 first_entered_rx.await.expect("first start must pause");
2012 let first_owner = state.client();
2013
2014 let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
2015 let second_state = state.clone();
2016 let second_starts = starts.clone();
2017 let second_stopped_owners = stopped_owners.clone();
2018 let second = tokio::spawn(async move {
2019 second_attempting_tx.send(()).unwrap();
2020 simulate_managed_session_start(
2021 second_state,
2022 "second-user".to_string(),
2023 None,
2024 second_starts,
2025 second_stopped_owners,
2026 )
2027 .await
2028 });
2029 second_attempting_rx.await.unwrap();
2030 tokio::task::yield_now().await;
2031 assert!(Arc::ptr_eq(&first_owner, &state.client()));
2032
2033 first_release_tx.send(()).unwrap();
2034 assert_eq!(
2035 first.await.expect("first start task"),
2036 SessionDisposition::Start
2037 );
2038 assert_eq!(
2039 second.await.expect("replacement start task"),
2040 SessionDisposition::Replace
2041 );
2042 assert_eq!(starts.load(Ordering::SeqCst), 2);
2043
2044 let active = state
2045 .managed_session
2046 .lock()
2047 .await
2048 .clone()
2049 .expect("active owner");
2050 assert_eq!(active.owner_epoch, 2);
2051 assert!(active.key.contains("second-user"));
2052 assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
2053 assert_eq!(
2054 stopped_owners.lock().await.as_slice(),
2055 [Arc::as_ptr(&first_owner) as usize]
2056 );
2057 assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
2058 }
2059
2060 #[test]
2061 fn derives_app_tag_from_api_key_by_default() {
2062 let config = OpenRtcTauriConfig {
2063 api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
2064 ..OpenRtcTauriConfig::default()
2065 };
2066
2067 assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
2068 }
2069
2070 #[tokio::test]
2071 async fn requested_native_transport_without_installer_keeps_base_route_available() {
2072 let state = OpenRtcTauriState::new(test_config());
2073 let client = state.client();
2074 ensure_requested_native_transports(
2075 &state,
2076 &client,
2077 Some(&requested_ble_config()),
2078 &native_transport_install_context(),
2079 )
2080 .await
2081 .expect("an unavailable optional transport must not fail base Iroh startup");
2082
2083 assert!(state.installed_native_transports.lock().await.is_empty());
2084 }
2085
2086 #[tokio::test]
2087 async fn native_transport_installer_is_idempotent_for_one_client() {
2088 let installs = Arc::new(AtomicUsize::new(0));
2089 let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
2090 Arc::new(FakeTransportInstaller {
2091 installs: installs.clone(),
2092 }),
2093 );
2094 let client = state.client();
2095 let config = requested_ble_config();
2096
2097 ensure_requested_native_transports(
2098 &state,
2099 &client,
2100 Some(&config),
2101 &native_transport_install_context(),
2102 )
2103 .await
2104 .expect("first install");
2105 ensure_requested_native_transports(
2106 &state,
2107 &client,
2108 Some(&config),
2109 &native_transport_install_context(),
2110 )
2111 .await
2112 .expect("idempotent install");
2113
2114 assert_eq!(installs.load(Ordering::SeqCst), 1);
2115 }
2116
2117 #[tokio::test]
2118 async fn unavailable_native_transport_keeps_base_route_available() {
2119 let state = OpenRtcTauriState::new(test_config())
2120 .with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
2121 let client = state.client();
2122
2123 ensure_requested_native_transports(
2124 &state,
2125 &client,
2126 Some(&requested_ble_config()),
2127 &native_transport_install_context(),
2128 )
2129 .await
2130 .expect("optional transport installation failure must not fail base Iroh startup");
2131
2132 assert!(state.installed_native_transports.lock().await.is_empty());
2133 }
2134
2135 #[test]
2136 fn invalid_or_missing_api_key_fails_closed() {
2137 assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
2138 let mut config = OpenRtcTauriConfig::default();
2139 config.api_key = "pk_test_short".to_string();
2140 assert!(config.validated_api_key().is_err());
2141 }
2142
2143 #[test]
2144 fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
2145 let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
2146 let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
2147 let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
2148 let api_key = format!("pk_test_{}", "c".repeat(40));
2149 std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
2150 std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
2151 std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
2152
2153 let config = OpenRtcTauriConfig::from_env();
2154
2155 assert_eq!(config.api_key, api_key);
2156 assert_eq!(
2157 config.app_tag().unwrap(),
2158 openrtc::app_tag_from_api_key(&api_key)
2159 );
2160
2161 match previous_api_key {
2162 Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
2163 None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
2164 }
2165 match previous_project {
2166 Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
2167 None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
2168 }
2169 match previous_app_tag {
2170 Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
2171 None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
2172 }
2173 }
2174
2175 #[test]
2176 fn native_state_has_one_constructor_owned_app_identity() {
2177 let config = test_config();
2178 let expected = openrtc::app_tag_from_api_key(&config.api_key);
2179 let state = OpenRtcTauriState::new(config);
2180 let first = state.client();
2181 let second = state.client();
2182
2183 assert_eq!(first.app_tag(), expected);
2184 assert!(Arc::ptr_eq(&first, &second));
2185 }
2186
2187 #[test]
2188 fn managed_session_uses_persisted_native_device_id() {
2189 let identity = openrtc::native_device::NativeDeviceIdentity {
2190 device_id: "persisted-native-device".to_string(),
2191 device_name: "Mac".to_string(),
2192 created_at_ms: 1,
2193 updated_at_ms: 1,
2194 name_source: None,
2195 system_info: None,
2196 };
2197
2198 assert_eq!(
2199 managed_session_device_id(Some(" desktop-e2e-native "), &identity),
2200 "persisted-native-device"
2201 );
2202 assert_eq!(
2203 managed_session_device_id(None, &identity),
2204 "persisted-native-device"
2205 );
2206 }
2207
2208 #[test]
2209 fn repeated_managed_session_triggers_converge_on_one_owner() {
2210 let key = "app:user:persisted-native-device";
2211 let cases = [
2212 ("react-remount", true, true, true, true),
2213 ("hmr", true, true, true, true),
2214 ("auth-refresh", true, true, true, true),
2215 ("resume", true, true, true, true),
2216 ("alias-change", true, true, true, true),
2217 ];
2218 for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
2219 assert_eq!(
2220 managed_session_disposition(
2221 Some(key),
2222 true,
2223 active_presence,
2224 active_auto,
2225 key,
2226 requested_presence,
2227 requested_auto,
2228 ),
2229 SessionDisposition::Reuse,
2230 "{label} must reuse the authoritative tuple"
2231 );
2232 }
2233
2234 assert_eq!(
2235 managed_session_disposition(Some(key), true, false, true, key, true, true),
2236 SessionDisposition::Refresh,
2237 "failed presence startup must retry idempotently"
2238 );
2239 assert_eq!(
2240 managed_session_disposition(
2241 Some(key),
2242 true,
2243 true,
2244 true,
2245 "app:other-user:persisted-native-device",
2246 true,
2247 true,
2248 ),
2249 SessionDisposition::Replace,
2250 "user switching must replace the previous lifecycle owner"
2251 );
2252 assert_eq!(
2253 managed_session_disposition(Some(key), false, true, true, key, true, true),
2254 SessionDisposition::Replace,
2255 "the same tuple on a different client epoch must replace the previous owner"
2256 );
2257 }
2258
2259 #[test]
2260 fn user_device_revocation_invalidates_the_cached_managed_ticket() {
2261 assert!(revokes_managed_session("user-device"));
2262 assert!(revokes_managed_session(" user-device "));
2263 assert!(!revokes_managed_session("share:example"));
2264 assert!(!revokes_managed_session(""));
2265 }
2266
2267 #[test]
2268 fn managed_session_metadata_uses_authoritative_device_id() {
2269 let raw = serde_json::json!({
2270 "deviceId": "persisted-native-device",
2271 "assistantDevice": {
2272 "deviceId": "desktop-e2e-native"
2273 }
2274 })
2275 .to_string();
2276
2277 let metadata =
2278 metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
2279 let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
2280
2281 assert_eq!(parsed["deviceId"], "desktop-e2e-native");
2282 assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
2283 }
2284}