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