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