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