use std::collections::{BTreeSet, HashMap};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use base64::Engine as _;
use ed25519_dalek::{Signature, VerifyingKey};
use openrtc::application_crypto_streams::{PeerRecvStream, PeerSendStream};
use serde::{Deserialize, Serialize};
#[cfg(feature = "managed-group-encryption")]
use sha2::{Digest, Sha256};
use tauri::ipc::{Channel, InvokeBody, Request, Response};
use tauri::{Manager, Runtime};
use tokio::sync::Mutex;
#[cfg(feature = "native-broadcast-moq")]
mod native_broadcast_moq;
const PLUGIN_NAME: &str = "openrtc-tauri-plugin";
const PEER_BI_STREAM_ID_HEADER: &str = "x-openrtc-stream-id";
const PEER_ID_HEADER: &str = "x-openrtc-peer-id";
const FANOUT_CAPABILITY_HEADER: &str = "x-openrtc-fanout-capability";
const FANOUT_SOURCE_PEER_HEADER: &str = "x-openrtc-fanout-source-peer";
const FANOUT_APP_TAG_HEADER: &str = "x-openrtc-fanout-app-tag";
const MEDIA_PUBLICATION_ID_HEADER: &str = "x-openrtc-publication-id";
const MEDIA_TIMESTAMP_US_HEADER: &str = "x-openrtc-timestamp-us";
const MEDIA_DURATION_US_HEADER: &str = "x-openrtc-duration-us";
const MEDIA_KEYFRAME_HEADER: &str = "x-openrtc-keyframe";
const MEDIA_DISCARDABLE_HEADER: &str = "x-openrtc-discardable";
fn request_body_bytes(request: &Request<'_>) -> Result<Vec<u8>, String> {
match request.body() {
InvokeBody::Raw(bytes) => Ok(bytes.clone()),
InvokeBody::Json(json) => serde_json::from_value::<Vec<u8>>(json.clone())
.map_err(|error| format!("invalid binary IPC payload: {error}")),
}
}
fn required_request_header(request: &Request<'_>, name: &str) -> Result<String, String> {
request
.headers()
.get(name)
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.ok_or_else(|| format!("{name} header is required"))
}
fn parsed_request_header<T>(request: &Request<'_>, name: &str) -> Result<T, String>
where
T: std::str::FromStr,
T::Err: std::fmt::Display,
{
required_request_header(request, name)?
.parse::<T>()
.map_err(|error| format!("invalid {name} header: {error}"))
}
pub type InstallFuture = std::pin::Pin<
Box<
dyn std::future::Future<Output = Result<Box<dyn std::any::Any + Send + Sync>, String>>
+ Send,
>,
>;
pub type BroadcastAdapterFuture<'a> = std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<Option<openrtc::broadcast::BroadcastAdapterObservation>, String>,
> + Send
+ 'a,
>,
>;
pub type BroadcastObjectsFuture<'a> =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<Vec<Vec<u8>>, String>> + Send + 'a>>;
pub trait NativeBroadcastAdapter: Send + Sync {
fn apply(
&self,
session: openrtc::broadcast::BroadcastSession,
action: openrtc::broadcast::BroadcastAdapterAction,
) -> BroadcastAdapterFuture<'_>;
fn take_inbound_objects(
&self,
_session: openrtc::broadcast::BroadcastSession,
_max: usize,
) -> BroadcastObjectsFuture<'_> {
Box::pin(async { Ok(Vec::new()) })
}
}
#[derive(Debug, Clone)]
pub struct InstallContext {
pub data_dir: PathBuf,
}
pub trait TransportInstaller: Send + Sync {
fn id(&self) -> &'static str;
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool;
fn install(
&self,
client: Arc<openrtc::client::Client>,
config: openrtc::client::TransportConfig,
context: InstallContext,
) -> InstallFuture;
}
pub trait DeviceKeySigner: Send + Sync {
fn public_jwk(&self, app_tag: &str) -> Result<serde_json::Value, String>;
fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String>;
fn offline_assurance(&self, _app_tag: &str) -> openrtc::offline::OfflineAssurance {
openrtc::offline::OfflineAssurance::Software
}
fn delete(&self, app_tag: &str) -> Result<(), String>;
fn read_secure_record(&self, _app_tag: &str, _key: &str) -> Result<Option<String>, String> {
Ok(None)
}
fn write_secure_record(&self, _app_tag: &str, _key: &str, _value: &str) -> Result<(), String> {
Err("OpenRTC 2.0 native certificate persistence requires a host secure store".to_string())
}
fn delete_secure_record(&self, _app_tag: &str, _key: &str) -> Result<(), String> {
Ok(())
}
}
struct OfflineDeviceSigner<'a> {
signer: &'a dyn DeviceKeySigner,
app_tag: &'a str,
}
impl openrtc::offline::OfflineSigner for OfflineDeviceSigner<'_> {
fn verifying_key(&self) -> anyhow::Result<VerifyingKey> {
let jwk = self
.signer
.public_jwk(self.app_tag)
.map_err(anyhow::Error::msg)?;
validate_public_device_jwk(&jwk).map_err(anyhow::Error::msg)?;
let encoded = jwk
.get("x")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| anyhow::anyhow!("OpenRTC device JWK is missing x"))?;
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(encoded)
.map_err(|error| anyhow::anyhow!("decode OpenRTC device JWK: {error}"))?;
let bytes: [u8; 32] = bytes
.try_into()
.map_err(|_| anyhow::anyhow!("OpenRTC device JWK must contain 32 key bytes"))?;
VerifyingKey::from_bytes(&bytes)
.map_err(|error| anyhow::anyhow!("parse OpenRTC device JWK: {error}"))
}
fn sign(&self, message: &[u8]) -> anyhow::Result<Signature> {
let bytes = self
.signer
.sign(self.app_tag, message)
.map_err(anyhow::Error::msg)?;
Signature::from_slice(&bytes)
.map_err(|error| anyhow::anyhow!("parse OpenRTC device signature: {error}"))
}
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct OfflineRuntimeSupport {
provisioning: bool,
local_mesh: bool,
cloud_required: bool,
#[serde(skip_serializing_if = "Option::is_none")]
reason: Option<&'static str>,
}
struct InstalledNativeTransport {
client_ptr: usize,
_runtime: Box<dyn std::any::Any + Send + Sync>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct OpenRtcTauriConfig {
pub api_key: String,
#[serde(default)]
pub data_dir: Option<PathBuf>,
#[serde(default)]
pub transport_config: Option<openrtc::client::TransportConfig>,
}
impl Default for OpenRtcTauriConfig {
fn default() -> Self {
Self {
api_key: String::new(),
data_dir: None,
transport_config: None,
}
}
}
impl OpenRtcTauriConfig {
pub fn from_env() -> Self {
let api_key = first_env(&[
"VITE_OPENRTC_KEY",
"VITE_OPENRTC_API_KEY",
"VITE_PLUTO_OPENRTC_API_KEY",
"OPENRTC_API_KEY",
]);
Self {
api_key: api_key.unwrap_or_default(),
data_dir: None,
transport_config: None,
}
}
pub fn validated_api_key(&self) -> Result<&str, String> {
openrtc::validate_api_key(&self.api_key)
.map_err(|error| format!("invalid OpenRTC 2.0 public API key: {error}"))
}
pub fn app_tag(&self) -> Result<String, String> {
self.validated_api_key().map(openrtc::app_tag_from_api_key)
}
}
fn first_env(names: &[&str]) -> Option<String> {
names.iter().find_map(|name| {
std::env::var(name)
.ok()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
})
}
#[derive(Default)]
struct TokenRelayState {
identity_credential: RwLock<Option<String>>,
}
impl TokenRelayState {
fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
let relay = self.clone();
Box::new(move || {
relay
.identity_credential
.read()
.ok()
.and_then(|guard| guard.clone())
})
}
fn set(&self, identity_credential: Option<String>) {
if let Ok(mut guard) = self.identity_credential.write() {
*guard = normalize_token(identity_credential);
}
}
}
fn normalize_token(value: Option<String>) -> Option<String> {
value
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
}
#[derive(Debug, Clone)]
struct NativeCapabilityRegistration {
avenue_kind: String,
avenue_id: String,
sparse_fanout: bool,
requested_architecture: Option<openrtc::native::RoomArchitectureMode>,
desired_revision: u64,
desired_peers: Vec<serde_json::Value>,
}
impl NativeCapabilityRegistration {
fn uses_sparse_fanout(&self) -> bool {
self.requested_architecture
.map(|requested| requested == openrtc::native::RoomArchitectureMode::Sparse)
.unwrap_or(self.sparse_fanout)
}
}
#[derive(Debug, Default)]
struct NativeCapabilityRegistry {
registrations: HashMap<String, NativeCapabilityRegistration>,
root_desired_revision: u64,
}
impl NativeCapabilityRegistry {
#[cfg(test)]
fn register(
&mut self,
capability_key: String,
avenue_kind: String,
avenue_id: String,
sparse_fanout: bool,
) -> Result<(), String> {
self.register_with_architecture(capability_key, avenue_kind, avenue_id, sparse_fanout, None)
}
fn register_with_architecture(
&mut self,
capability_key: String,
avenue_kind: String,
avenue_id: String,
sparse_fanout: bool,
requested_architecture: Option<openrtc::native::RoomArchitectureMode>,
) -> Result<(), String> {
if requested_architecture.is_some() && avenue_kind != "room" {
return Err("native room architecture is valid only for room avenues".to_string());
}
if let Some(existing) = self.registrations.get(&capability_key) {
if existing.avenue_kind == avenue_kind
&& existing.avenue_id == avenue_id
&& existing.sparse_fanout == sparse_fanout
&& existing.requested_architecture == requested_architecture
{
return Ok(());
}
return Err(format!(
"native capability key {capability_key} is already registered for another avenue"
));
}
self.registrations.insert(
capability_key,
NativeCapabilityRegistration {
avenue_kind,
avenue_id,
sparse_fanout,
requested_architecture,
desired_revision: 0,
desired_peers: Vec::new(),
},
);
Ok(())
}
fn capability_keys_for_identity(
&self,
connection_id: Option<&str>,
device_id: Option<&str>,
device_id_hint: Option<&str>,
remote_node_id: Option<&str>,
) -> Vec<String> {
let identities = [connection_id, device_id, device_id_hint, remote_node_id]
.into_iter()
.flatten()
.map(str::trim)
.filter(|value| !value.is_empty())
.collect::<BTreeSet<_>>();
self.registrations
.iter()
.filter_map(|(key, registration)| {
registration
.desired_peers
.iter()
.any(|peer| {
["connectionId", "deviceId", "nodeId"]
.into_iter()
.filter_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
.map(str::trim)
.any(|value| identities.contains(value))
})
.then(|| key.clone())
})
.collect::<BTreeSet<_>>()
.into_iter()
.collect()
}
fn capability_keys_for_state(&self, snapshot: &openrtc::client::StateSnapshot) -> Vec<String> {
self.capability_keys_for_identity(
Some(&snapshot.connection_id),
snapshot.device_id.as_deref(),
snapshot.device_id_hint.as_deref(),
snapshot.remote_node_id.as_deref(),
)
}
fn capability_keys_for_peer_data(
&self,
event: &openrtc::client::NativePeerDataEvent,
) -> Vec<String> {
if let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) {
if let Some(capability) = value
.get("capability")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
{
if self.registrations.contains_key(capability) {
return vec![capability.to_string()];
}
return Vec::new();
}
}
Vec::new()
}
fn projected_stream_capability(
&self,
channel: &openrtc::stream_metadata::ChannelMetadata,
) -> Option<String> {
let explicit = channel
.metadata
.as_ref()
.and_then(|metadata| metadata.get("openrtcCapability"))
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty());
match explicit {
Some(key) if self.registrations.contains_key(key) => Some(key.to_string()),
_ => None,
}
}
fn aggregate_desired_peers(&mut self) -> Result<(u64, String), String> {
let role_priority = |value: &serde_json::Value| match value
.get("topologyRole")
.and_then(serde_json::Value::as_str)
{
None => 3_u8,
Some("active") => 2,
Some("backup") => 1,
Some(_) => 0,
};
let route_evidence = |value: &serde_json::Value| {
["ticket", "nodeId"]
.into_iter()
.filter(|field| {
value
.get(field)
.and_then(serde_json::Value::as_str)
.is_some_and(|part| !part.trim().is_empty())
})
.count()
};
let mut references =
std::collections::BTreeMap::<String, Vec<(&str, &serde_json::Value)>>::new();
let mut registrations = self.registrations.iter().collect::<Vec<_>>();
registrations.sort_by(|(left, _), (right, _)| left.cmp(right));
for (capability_key, registration) in registrations {
for peer in ®istration.desired_peers {
let identity = ["deviceId", "nodeId", "ticket"]
.into_iter()
.find_map(|field| peer.get(field).and_then(serde_json::Value::as_str))
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| "desired peer has no stable identity".to_string())?;
references
.entry(identity.to_string())
.or_default()
.push((capability_key.as_str(), peer));
}
}
self.root_desired_revision = self.root_desired_revision.saturating_add(1);
let peers = references
.into_values()
.map(|refs| {
let semantic_priority = refs
.iter()
.map(|(_, peer)| role_priority(peer))
.max()
.unwrap_or_default();
let mut evidence = refs[0];
for candidate in refs.iter().copied().skip(1) {
if (route_evidence(candidate.1), role_priority(candidate.1))
> (route_evidence(evidence.1), role_priority(evidence.1))
{
evidence = candidate;
}
}
let mut peer = evidence.1.clone();
if semantic_priority == 3 {
if let Some(object) = peer.as_object_mut() {
object.remove("topologyRole");
object.remove("topologyRevision");
}
} else {
peer["topologyRole"] = serde_json::Value::String(
if semantic_priority == 2 {
"active"
} else {
"backup"
}
.to_string(),
);
peer["topologyRevision"] = serde_json::Value::from(self.root_desired_revision);
}
peer
})
.collect::<Vec<_>>();
let payload = serde_json::to_string(&peers)
.map_err(|error| format!("serialize aggregated desired peers: {error}"))?;
Ok((self.root_desired_revision, payload))
}
}
fn required_native_capability_part(value: String, label: &str) -> Result<String, String> {
let value = value.trim().to_string();
if value.is_empty()
|| value.len() > 192
|| !value
.chars()
.all(|character| character.is_ascii_alphanumeric() || "_.:@-".contains(character))
{
return Err(format!("native {label} is invalid"));
}
Ok(value)
}
pub struct OpenRtcTauriState {
client: RwLock<Arc<openrtc::client::Client>>,
config: OpenRtcTauriConfig,
token_relay: Arc<TokenRelayState>,
data_dir: Option<PathBuf>,
connection_state_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
peer_data_forwarder: std::sync::Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
presence_loop_active: Mutex<bool>,
subscriptions: Mutex<HashMap<String, tokio::task::JoinHandle<()>>>,
native_projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
projected_stream_offered: AtomicU64,
projected_stream_decoded: AtomicU64,
projected_stream_projected: AtomicU64,
projected_stream_unhandled_non_channel: AtomicU64,
projected_stream_unhandled_other_channel: AtomicU64,
projected_stream_unhandled_no_subscription: AtomicU64,
projected_stream_failures: AtomicU64,
projected_stream_last_transport_stable_id: AtomicU64,
projected_stream_validation_checks: AtomicU64,
projected_stream_validation_rejections: AtomicU64,
projected_stream_last_validated_transport_stable_id: AtomicU64,
projected_stream_authorized: AtomicU64,
projected_stream_unauthorized: AtomicU64,
projected_stream_last_channel: Mutex<Option<String>>,
projected_stream_last_protocol: Mutex<Option<String>>,
projected_stream_last_connection_id: Mutex<Option<String>>,
projected_stream_last_remote_node_id: Mutex<Option<String>>,
peer_bi_streams: Mutex<HashMap<String, PeerBiStreamHandle>>,
peer_uni_streams: Mutex<HashMap<String, Arc<Mutex<Option<PeerSendStream>>>>>,
portable_media: Mutex<openrtc::media::PortableMediaSession>,
broadcast_sessions: Mutex<HashMap<String, openrtc::broadcast::BroadcastSession>>,
broadcast_signers: Mutex<HashMap<String, openrtc::broadcast::BroadcastPublisherSigner>>,
broadcast_drivers: Mutex<HashMap<String, Arc<Mutex<()>>>>,
managed_session_start_guard: Mutex<()>,
managed_session_next_owner_epoch: AtomicU64,
managed_session: Mutex<Option<ManagedSessionRecord>>,
capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
native_transport_installers: Vec<Arc<dyn TransportInstaller>>,
installed_native_transports: Mutex<HashMap<&'static str, InstalledNativeTransport>>,
native_device_key_signer: Option<Arc<dyn DeviceKeySigner>>,
native_broadcast_adapter: Option<Arc<dyn NativeBroadcastAdapter>>,
#[cfg(feature = "managed-group-encryption")]
managed_groups: Mutex<HashMap<String, openrtc::native::NativeManagedGroupController>>,
#[cfg(feature = "native-broadcast-moq")]
native_broadcast_moq_adapter: Option<Arc<native_broadcast_moq::NativeBroadcastMoqAdapter>>,
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeRuntimeStatusResponse {
#[serde(flatten)]
status: openrtc::client::RuntimeStatus,
product_maturity: openrtc::client::ProductCapabilityMaturity,
}
impl OpenRtcTauriState {
pub fn new(config: OpenRtcTauriConfig) -> Self {
openrtc::ensure_rustls();
let token_relay = Arc::new(TokenRelayState::default());
let client = build_client(&config, &token_relay)
.expect("OpenRTC Tauri 2.0 requires a valid public API key");
let data_dir = config.data_dir.clone();
#[cfg(feature = "native-broadcast-moq")]
let native_broadcast_moq_adapter =
Arc::new(native_broadcast_moq::NativeBroadcastMoqAdapter::new());
Self {
client: RwLock::new(client),
config,
token_relay,
data_dir,
connection_state_forwarder: std::sync::Mutex::new(None),
peer_data_forwarder: std::sync::Mutex::new(None),
presence_loop_active: Mutex::new(false),
subscriptions: Mutex::new(HashMap::new()),
native_projection_subscription: Arc::new(Mutex::new(None)),
projected_stream_offered: AtomicU64::new(0),
projected_stream_decoded: AtomicU64::new(0),
projected_stream_projected: AtomicU64::new(0),
projected_stream_unhandled_non_channel: AtomicU64::new(0),
projected_stream_unhandled_other_channel: AtomicU64::new(0),
projected_stream_unhandled_no_subscription: AtomicU64::new(0),
projected_stream_failures: AtomicU64::new(0),
projected_stream_last_transport_stable_id: AtomicU64::new(0),
projected_stream_validation_checks: AtomicU64::new(0),
projected_stream_validation_rejections: AtomicU64::new(0),
projected_stream_last_validated_transport_stable_id: AtomicU64::new(0),
projected_stream_authorized: AtomicU64::new(0),
projected_stream_unauthorized: AtomicU64::new(0),
projected_stream_last_channel: Mutex::new(None),
projected_stream_last_protocol: Mutex::new(None),
projected_stream_last_connection_id: Mutex::new(None),
projected_stream_last_remote_node_id: Mutex::new(None),
peer_bi_streams: Mutex::new(HashMap::new()),
peer_uni_streams: Mutex::new(HashMap::new()),
portable_media: Mutex::new(openrtc::media::PortableMediaSession::default()),
broadcast_sessions: Mutex::new(HashMap::new()),
broadcast_signers: Mutex::new(HashMap::new()),
broadcast_drivers: Mutex::new(HashMap::new()),
managed_session_start_guard: Mutex::new(()),
managed_session_next_owner_epoch: AtomicU64::new(0),
managed_session: Mutex::new(None),
capability_registry: Arc::new(Mutex::new(NativeCapabilityRegistry::default())),
native_transport_installers: Vec::new(),
installed_native_transports: Mutex::new(HashMap::new()),
native_device_key_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_groups: Mutex::new(HashMap::new()),
native_broadcast_adapter: {
#[cfg(feature = "native-broadcast-moq")]
{
Some(native_broadcast_moq_adapter.clone())
}
#[cfg(not(feature = "native-broadcast-moq"))]
{
None
}
},
#[cfg(feature = "native-broadcast-moq")]
native_broadcast_moq_adapter: Some(native_broadcast_moq_adapter),
}
}
pub fn with_native_transport_installer(
mut self,
installer: Arc<dyn TransportInstaller>,
) -> Self {
self.native_transport_installers.push(installer);
self
}
pub fn with_native_device_key_signer(mut self, signer: Arc<dyn DeviceKeySigner>) -> Self {
self.native_device_key_signer = Some(signer);
self
}
fn offline_runtime_support(&self) -> OfflineRuntimeSupport {
let provisioning = self.native_device_key_signer.is_some();
OfflineRuntimeSupport {
provisioning,
local_mesh: cfg!(feature = "transport-lan"),
cloud_required: false,
reason: if !provisioning {
Some("native host device signer is unavailable")
} else {
(!cfg!(feature = "transport-lan"))
.then_some("native host was built without transport-lan")
},
}
}
fn project_public_runtime_status(
&self,
status: openrtc::client::RuntimeStatus,
) -> NativeRuntimeStatusResponse {
use openrtc::client::CapabilityMaturity;
let mut product_maturity = self.client().product_capability_maturity();
product_maturity.offline_edge = if self.native_device_key_signer.is_none() {
CapabilityMaturity::Unavailable
} else if cfg!(feature = "transport-lan") {
CapabilityMaturity::Preview
} else {
CapabilityMaturity::SupportOnly
};
product_maturity.broadcast = CapabilityMaturity::Unavailable;
NativeRuntimeStatusResponse {
status,
product_maturity,
}
}
pub fn with_native_broadcast_adapter(
mut self,
adapter: Arc<dyn NativeBroadcastAdapter>,
) -> Self {
self.native_broadcast_adapter = Some(adapter);
#[cfg(feature = "native-broadcast-moq")]
{
self.native_broadcast_moq_adapter = None;
}
self
}
#[cfg(feature = "native-broadcast-moq")]
async fn resolve_native_broadcast_access(
&self,
grant: &str,
publication_verifying_key: Option<&[u8; 32]>,
expected_issuer_public_key: &[u8; 32],
) -> Result<native_broadcast_moq::NativeBroadcastRelayAccess, String> {
let api_key = self.config.validated_api_key()?;
let app_tag = self.client().app_tag().to_string();
let nonce = uuid::Uuid::new_v4().simple().to_string();
let issued_at = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let (challenge, publication_key) = native_broadcast_moq::broadcast_access_challenge(
api_key,
grant,
publication_verifying_key,
&nonce,
issued_at,
);
let device_proof = native_broadcast_moq::NativeBroadcastDeviceProof {
public_key_jwk: self.device_public_key(&app_tag)?,
signature: self.sign_device_proof(&app_tag, &challenge)?,
nonce,
issued_at,
};
native_broadcast_moq::resolve_broadcast_access(
openrtc::native::OPENRTC_PRODUCTION_CONTROL_PLANE,
api_key,
grant,
publication_key.as_deref(),
device_proof,
expected_issuer_public_key,
)
.await
}
fn device_public_key(&self, app_tag: &str) -> Result<serde_json::Value, String> {
let app_tag = required_app_tag(app_tag)?;
let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
"OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
})?;
let value = signer.public_jwk(app_tag)?;
validate_public_device_jwk(&value)?;
Ok(value)
}
fn sign_device_proof(&self, app_tag: &str, challenge: &str) -> Result<String, String> {
let app_tag = required_app_tag(app_tag)?;
if challenge.is_empty() || challenge.len() > 2_048 {
return Err("OpenRTC device challenge is invalid".to_string());
}
let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
"OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
})?;
let signature = signer.sign(app_tag, challenge.as_bytes())?;
if signature.len() != 64 {
return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
}
Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
}
fn sign_device_message(&self, app_tag: &str, message: &[u8]) -> Result<String, String> {
let app_tag = required_app_tag(app_tag)?;
if message.is_empty() || message.len() > 96 * 1024 {
return Err("OpenRTC device message is invalid".to_string());
}
let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
"OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
})?;
let signature = signer.sign(app_tag, message)?;
if signature.len() != 64 {
return Err("OpenRTC device signer returned an invalid Ed25519 signature".to_string());
}
Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature))
}
fn delete_device_key(&self, app_tag: &str) -> Result<(), String> {
let app_tag = required_app_tag(app_tag)?;
let signer = self.native_device_key_signer.as_ref().ok_or_else(|| {
"OpenRTC 2.0 native device proof requires a host secure-storage signer".to_string()
})?;
signer.delete(app_tag)
}
fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
let app_tag = required_app_tag(app_tag)?;
let key = required_secure_record_key(key)?;
self.native_device_key_signer
.as_ref()
.ok_or_else(|| {
"OpenRTC 2.0 native certificate persistence requires a host secure store"
.to_string()
})?
.read_secure_record(app_tag, key)
}
fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
let app_tag = required_app_tag(app_tag)?;
let key = required_secure_record_key(key)?;
if value.is_empty() || value.len() > 16 * 1024 {
return Err("OpenRTC secure record is invalid".to_string());
}
self.native_device_key_signer
.as_ref()
.ok_or_else(|| {
"OpenRTC 2.0 native certificate persistence requires a host secure store"
.to_string()
})?
.write_secure_record(app_tag, key, value)
}
fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
let app_tag = required_app_tag(app_tag)?;
let key = required_secure_record_key(key)?;
self.native_device_key_signer
.as_ref()
.ok_or_else(|| {
"OpenRTC 2.0 native certificate persistence requires a host secure store"
.to_string()
})?
.delete_secure_record(app_tag, key)
}
pub fn client(&self) -> Arc<openrtc::client::Client> {
self.client
.read()
.map(|guard| guard.clone())
.expect("OpenRTC Tauri client state is poisoned")
}
pub fn set_identity_credential(&self, identity_credential: Option<String>) {
self.token_relay.set(identity_credential);
}
fn allocate_managed_session_owner_epoch(&self) -> u64 {
self.managed_session_next_owner_epoch
.fetch_add(1, Ordering::SeqCst)
+ 1
}
fn replace_connection_state_forwarder(&self, client: Arc<openrtc::client::Client>) {
let capability_registry = self.capability_registry.clone();
let projection_subscription = self.native_projection_subscription.clone();
let next = tauri::async_runtime::spawn(async move {
forward_connection_state_events(client, capability_registry, projection_subscription)
.await;
});
if let Ok(mut guard) = self.connection_state_forwarder.lock() {
if let Some(previous) = guard.replace(next) {
previous.abort();
}
}
}
fn replace_peer_data_forwarder(&self, client: Arc<openrtc::client::Client>) {
let capability_registry = self.capability_registry.clone();
let projection_subscription = self.native_projection_subscription.clone();
let next = tauri::async_runtime::spawn(async move {
forward_peer_data_events(client, capability_registry, projection_subscription).await;
});
if let Ok(mut guard) = self.peer_data_forwarder.lock() {
if let Some(previous) = guard.replace(next) {
previous.abort();
}
}
}
}
fn build_client(
config: &OpenRtcTauriConfig,
token_relay: &Arc<TokenRelayState>,
) -> Result<Arc<openrtc::client::Client>, String> {
let mut builder = openrtc::client::Client::builder(
config.validated_api_key()?.to_string(),
token_relay.token_provider(),
)
.map_err(|error| error.to_string())?;
if let Some(transport_config) = config.transport_config.clone() {
builder = builder.transport_config(transport_config);
}
Ok(Arc::new(builder.build()))
}
fn requested_local_device_id(value: Option<&str>) -> Option<String> {
value
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn managed_session_device_id(
_requested: Option<&str>,
native_identity: &openrtc::native_device::NativeDeviceIdentity,
) -> String {
native_identity.device_id.clone()
}
#[cfg(test)]
fn metadata_with_authoritative_device_id(
metadata: Option<String>,
device_id: &str,
) -> Option<String> {
let device_id = device_id.trim();
if device_id.is_empty() {
return metadata;
}
let Some(raw_metadata) = metadata else {
return Some(serde_json::json!({ "deviceId": device_id }).to_string());
};
match serde_json::from_str::<serde_json::Value>(&raw_metadata) {
Ok(serde_json::Value::Object(mut map)) => {
map.insert(
"deviceId".to_string(),
serde_json::Value::String(device_id.to_string()),
);
Some(serde_json::Value::Object(map).to_string())
}
_ => Some(
serde_json::json!({
"deviceId": device_id,
"metadata": raw_metadata,
})
.to_string(),
),
}
}
struct PeerBiStreamHandle {
send: Arc<Mutex<Option<PeerSendStream>>>,
recv: Option<PeerRecvStream>,
read_task: Option<tokio::task::JoinHandle<()>>,
}
struct NativeProjectionSubscription {
request_id: String,
channel: Channel<NativeProjectionEvent>,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct NativeProjectionEvent {
request_id: String,
kind: &'static str,
capability_keys: Vec<String>,
payload: serde_json::Value,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
pub struct OpenBiResult {
stream_id: String,
connection_id: Option<String>,
remote_node_id: String,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
pub struct OpenUniResult {
stream_id: String,
connection_id: Option<String>,
remote_node_id: String,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct IncomingPeerBiStreamEvent {
request_id: String,
stream_id: String,
connection_id: Option<String>,
remote_node_id: String,
transport_stable_id: u64,
channel: Option<openrtc::stream_metadata::ChannelMetadata>,
application_authorized: bool,
capability_key: String,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct ProjectedStreamDiagnostics {
subscription_active: bool,
offered: u64,
decoded: u64,
projected: u64,
unhandled_non_channel: u64,
unhandled_other_channel: u64,
unhandled_no_subscription: u64,
failures: u64,
last_transport_stable_id: u64,
validation_checks: u64,
validation_rejections: u64,
last_validated_transport_stable_id: u64,
authorized: u64,
unauthorized: u64,
last_channel: Option<String>,
last_protocol: Option<String>,
last_connection_id: Option<String>,
last_remote_node_id: Option<String>,
}
pub enum ProjectedPeerBiHandoff {
Handled,
Unhandled {
send: PeerSendStream,
recv: PeerRecvStream,
},
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct PeerBiStreamClosedEvent {
r#type: &'static str,
error: Option<String>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct StartSessionResult {
local_node_id: String,
ticket_scope: Option<String>,
ticket: Option<String>,
presence_started: bool,
auto_connect_started: bool,
local_device: openrtc::native_device::NativeDeviceIdentity,
}
#[derive(Clone)]
struct ManagedSessionRecord {
key: String,
owner_client: Arc<openrtc::client::Client>,
owner_epoch: u64,
result: StartSessionResult,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SessionDisposition {
Start,
Reuse,
Refresh,
Replace,
}
fn managed_session_disposition(
active_key: Option<&str>,
active_client_matches: bool,
active_presence_started: bool,
active_auto_connect_started: bool,
requested_key: &str,
requested_presence: bool,
requested_auto_connect: bool,
) -> SessionDisposition {
let Some(active_key) = active_key else {
return SessionDisposition::Start;
};
if !active_client_matches || active_key != requested_key {
return SessionDisposition::Replace;
}
if (!requested_presence || active_presence_started)
&& (!requested_auto_connect || active_auto_connect_started)
{
SessionDisposition::Reuse
} else {
SessionDisposition::Refresh
}
}
fn revokes_managed_session(scope: &str) -> bool {
scope.trim() == "user-device"
}
pub fn init<R: Runtime>(config: OpenRtcTauriConfig) -> tauri::plugin::TauriPlugin<R> {
init_with_state(OpenRtcTauriState::new(config))
}
pub fn init_with_state<R: Runtime>(state: OpenRtcTauriState) -> tauri::plugin::TauriPlugin<R> {
tauri::plugin::Builder::new(PLUGIN_NAME)
.setup(move |app, _api| {
let client = state.client();
let transport_config = state.config.transport_config.clone();
let install_context = InstallContext {
data_dir: app_data_dir(app, &state).map_err(std::io::Error::other)?,
};
tauri::async_runtime::block_on(ensure_requested_native_transports(
&state,
&client,
transport_config.as_ref(),
&install_context,
))
.map_err(std::io::Error::other)?;
state.replace_connection_state_forwarder(client.clone());
state.replace_peer_data_forwarder(client);
app.manage(state);
Ok(())
})
.invoke_handler(tauri::generate_handler![
openrtc_set_identity_credential,
openrtc_device_public_key,
openrtc_sign_device_proof,
openrtc_delete_device_key,
openrtc_read_secure_record,
openrtc_write_secure_record,
openrtc_delete_secure_record,
openrtc_offline_runtime_support,
openrtc_create_offline_enrollment_request,
openrtc_verify_offline_enrollment_request,
rtc_native_status,
get_rtc_local_device_info,
update_rtc_local_device_name,
get_iroh_node_id,
start_iroh_node,
get_iroh_endpoint_ticket,
register_session_token,
get_endpoint_ticket_with_token,
validate_session_token,
revoke_session_tokens_by_scope,
stop_rtc_presence_loop,
register_rtc_capability,
openrtc_init_managed_room_group,
openrtc_handle_managed_room_prepare_page,
openrtc_handle_managed_room_artifact_chunk,
openrtc_seal_managed_room_payload,
openrtc_open_managed_room_payload,
openrtc_forget_managed_room_group,
unregister_rtc_capability,
start_rtc_managed_session,
start_rtc_external_auto_connect,
submit_rtc_desired_peers,
stop_rtc_auto_connect,
notify_rtc_network_change,
set_rtc_transport_priority,
connect_to_device,
disconnect_device,
set_auto_connect_excluded,
set_rtc_external_auto_connect_excluded,
resolve_rtc_peer_connection_records,
resolve_rtc_peer_identity,
get_rtc_peer_session,
list_rtc_peer_sessions,
list_rtc_managed_connections,
wait_for_rtc_settled_peer,
list_rtc_connection_states,
get_rtc_connection_state,
stop_rtc_subscription,
start_native_projection,
get_projected_stream_diagnostics,
is_current_transport_stable_id,
send_peer_message,
encode_sparse_fanout_message,
accept_sparse_fanout_message,
sparse_fanout_diagnostics,
prepare_openrtc_broadcast_publisher,
release_openrtc_broadcast_publisher,
open_openrtc_broadcast,
begin_openrtc_broadcast_publication,
publish_openrtc_broadcast_sample,
receive_openrtc_broadcast_media,
revoke_openrtc_broadcast,
close_openrtc_broadcast,
record_sparse_fanout_forward_queue_drop,
is_peer_connected,
open_peer_bi_stream,
open_peer_bi_transport_only_stream,
open_peer_uni_stream,
write_peer_bi_stream,
start_peer_bi_stream_read,
cancel_peer_bi_stream_read,
close_peer_bi_stream,
write_peer_uni_stream,
close_peer_uni_stream,
begin_openrtc_media_publication,
encode_openrtc_media_sample,
pause_openrtc_media_publication,
retire_openrtc_media_publication,
retire_openrtc_media_receiver,
decode_openrtc_media_chunk,
decode_openrtc_media_control,
])
.build()
}
pub fn init_from_env<R: Runtime>() -> tauri::plugin::TauriPlugin<R> {
init(OpenRtcTauriConfig::from_env())
}
async fn send_native_projection<T: Serialize>(
subscription: &Arc<Mutex<Option<NativeProjectionSubscription>>>,
kind: &'static str,
capability_keys: Vec<String>,
payload: &T,
) {
if capability_keys.is_empty() {
return;
}
let Ok(payload) = serde_json::to_value(payload) else {
return;
};
let target = subscription.lock().await.as_ref().map(|subscription| {
(
subscription.request_id.clone(),
subscription.channel.clone(),
)
});
let Some((request_id, channel)) = target else {
return;
};
let _ = channel.send(NativeProjectionEvent {
request_id,
kind,
capability_keys,
payload,
});
}
async fn forward_connection_state_events(
client: Arc<openrtc::client::Client>,
capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
) {
let mut rx = client.connection_state_updates();
while let Ok(snapshot) = rx.recv().await {
let keys = capability_registry
.lock()
.await
.capability_keys_for_state(&snapshot);
send_native_projection(
&projection_subscription,
"connection-state",
keys,
&snapshot,
)
.await;
}
}
async fn forward_peer_data_events(
client: Arc<openrtc::client::Client>,
capability_registry: Arc<Mutex<NativeCapabilityRegistry>>,
projection_subscription: Arc<Mutex<Option<NativeProjectionSubscription>>>,
) {
let mut rx = client.subscribe_native_peer_data();
while let Ok(event) = rx.recv().await {
let keys = capability_registry
.lock()
.await
.capability_keys_for_peer_data(&event);
send_native_projection(&projection_subscription, "peer-data", keys, &event).await;
}
}
async fn replay_current_connection_states(state: &OpenRtcTauriState) {
let snapshots = state.client().connection_states().await;
let registry = state.capability_registry.lock().await;
for snapshot in snapshots {
let keys = registry.capability_keys_for_state(&snapshot);
send_native_projection(
&state.native_projection_subscription,
"connection-state",
keys,
&snapshot,
)
.await;
}
}
fn app_data_dir<R: Runtime>(
app: &tauri::AppHandle<R>,
state: &OpenRtcTauriState,
) -> Result<PathBuf, String> {
state
.data_dir
.clone()
.or_else(|| app.path().app_data_dir().ok())
.ok_or_else(|| "failed to resolve OpenRTC app data directory".to_string())
}
async fn ensure_iroh_node(client: Arc<openrtc::client::Client>) -> Result<String, String> {
if let Some(node_id) = client.current_node_id().await {
return Ok(node_id);
}
client
.init_iroh(None, Vec::new())
.await
.map_err(|error| format!("failed to initialize OpenRTC Iroh node: {error}"))
}
async fn ensure_requested_native_transports(
state: &OpenRtcTauriState,
client: &Arc<openrtc::client::Client>,
transports: Option<&openrtc::client::TransportConfig>,
context: &InstallContext,
) -> Result<(), String> {
let Some(transports) = transports else {
return Ok(());
};
let ble_requested = transports.ble.as_ref().is_some_and(|config| config.enabled);
if ble_requested
&& !state
.native_transport_installers
.iter()
.any(|installer| installer.is_requested(transports))
{
eprintln!(
"[openrtc-tauri][transport] BLE requested but unavailable: this host did not register a BLE transport installer; continuing on the Iroh base route"
);
return Ok(());
}
let client_ptr = Arc::as_ptr(client) as usize;
for installer in &state.native_transport_installers {
if !installer.is_requested(transports) {
continue;
}
let already_installed = state
.installed_native_transports
.lock()
.await
.get(installer.id())
.is_some_and(|installed| installed.client_ptr == client_ptr);
if already_installed {
continue;
}
if client.current_node_id().await.is_some() {
eprintln!(
"[openrtc-tauri][transport] {} requested after native node startup; continuing on the Iroh base route",
installer.id()
);
continue;
}
let runtime = match installer
.install(client.clone(), transports.clone(), context.clone())
.await
{
Ok(runtime) => runtime,
Err(error) => {
eprintln!(
"[openrtc-tauri][transport] {} unavailable: {}; continuing on the Iroh base route",
installer.id(),
error
);
continue;
}
};
state.installed_native_transports.lock().await.insert(
installer.id(),
InstalledNativeTransport {
client_ptr,
_runtime: runtime,
},
);
}
Ok(())
}
async fn register_peer_bi_stream(
state: &OpenRtcTauriState,
connection_id: Option<String>,
remote_node_id: String,
send: PeerSendStream,
recv: PeerRecvStream,
) -> OpenBiResult {
let stream_id = uuid::Uuid::new_v4().to_string();
let send = Arc::new(Mutex::new(Some(send)));
state.peer_bi_streams.lock().await.insert(
stream_id.clone(),
PeerBiStreamHandle {
send,
recv: Some(recv),
read_task: None,
},
);
OpenBiResult {
stream_id,
connection_id,
remote_node_id,
}
}
async fn read_peer_exact(recv: &mut PeerRecvStream, buffer: &mut [u8]) -> std::io::Result<()> {
let mut offset = 0;
while offset < buffer.len() {
let read = recv.read(&mut buffer[offset..]).await?;
if read == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"peer stream ended during channel classification",
));
}
offset += read;
}
Ok(())
}
fn is_openrtc_projected_channel(channel: &openrtc::stream_metadata::ChannelMetadata) -> bool {
channel.channel_id == openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID
}
pub async fn handoff_projected_peer_bi_stream<R: Runtime>(
app: tauri::AppHandle<R>,
connection_id: Option<String>,
remote_node_id: String,
transport_stable_id: u64,
send: PeerSendStream,
mut recv: PeerRecvStream,
) -> Result<ProjectedPeerBiHandoff, String> {
let state = app.state::<OpenRtcTauriState>();
state
.projected_stream_offered
.fetch_add(1, Ordering::Relaxed);
state
.projected_stream_last_transport_stable_id
.store(transport_stable_id, Ordering::Relaxed);
*state.projected_stream_last_connection_id.lock().await = connection_id.clone();
*state.projected_stream_last_remote_node_id.lock().await = Some(remote_node_id.clone());
let mut envelope = Vec::new();
let channel = loop {
match openrtc::stream_metadata::decode_prefix(&envelope) {
openrtc::stream_metadata::DecodeDecision::NotMatched => {
state
.projected_stream_unhandled_non_channel
.fetch_add(1, Ordering::Relaxed);
return Ok(ProjectedPeerBiHandoff::Unhandled {
send,
recv: recv.with_plaintext_prefix(envelope),
});
}
openrtc::stream_metadata::DecodeDecision::Decoded(prefix) => break prefix.channel,
openrtc::stream_metadata::DecodeDecision::NeedMore(required) => {
let missing = required.saturating_sub(envelope.len());
if missing == 0 {
return Err("channel-envelope decoder made no progress".to_string());
}
let mut bytes = vec![0u8; missing];
if let Err(error) = read_peer_exact(&mut recv, &mut bytes).await {
state
.projected_stream_failures
.fetch_add(1, Ordering::Relaxed);
return Err(format!("read projected channel envelope: {error}"));
}
envelope.extend_from_slice(&bytes);
}
}
};
state
.projected_stream_decoded
.fetch_add(1, Ordering::Relaxed);
*state.projected_stream_last_channel.lock().await = Some(channel.channel_id.clone());
*state.projected_stream_last_protocol.lock().await = channel
.metadata
.as_ref()
.and_then(|metadata| metadata.get("protocol"))
.and_then(serde_json::Value::as_str)
.map(str::to_string);
if !is_openrtc_projected_channel(&channel) {
state
.projected_stream_unhandled_other_channel
.fetch_add(1, Ordering::Relaxed);
return Ok(ProjectedPeerBiHandoff::Unhandled {
send,
recv: recv.with_plaintext_prefix(envelope),
});
}
let Some(capability_key) = state
.capability_registry
.lock()
.await
.projected_stream_capability(&channel)
else {
state
.projected_stream_unhandled_other_channel
.fetch_add(1, Ordering::Relaxed);
return Ok(ProjectedPeerBiHandoff::Unhandled {
send,
recv: recv.with_plaintext_prefix(envelope),
});
};
let subscription = state.native_projection_subscription.lock().await;
let Some((request_id, projection_channel)) = subscription.as_ref().map(|subscription| {
(
subscription.request_id.clone(),
subscription.channel.clone(),
)
}) else {
state
.projected_stream_unhandled_no_subscription
.fetch_add(1, Ordering::Relaxed);
return Ok(ProjectedPeerBiHandoff::Unhandled {
send,
recv: recv.with_plaintext_prefix(envelope),
});
};
drop(subscription);
let application_authorized = connection_id.as_deref().is_some_and(|connection_id| {
matches!(
state.client().session_admission(connection_id),
openrtc::session_token::SessionAdmission::Accepted { .. }
)
});
if application_authorized {
state
.projected_stream_authorized
.fetch_add(1, Ordering::Relaxed);
} else {
state
.projected_stream_unauthorized
.fetch_add(1, Ordering::Relaxed);
}
let result = register_peer_bi_stream(
state.inner(),
connection_id,
remote_node_id.clone(),
send,
recv,
)
.await;
let event = IncomingPeerBiStreamEvent {
request_id: request_id.clone(),
stream_id: result.stream_id.clone(),
connection_id: result.connection_id,
remote_node_id,
transport_stable_id,
channel: Some(channel),
application_authorized,
capability_key: capability_key.clone(),
};
let payload = serde_json::to_value(event)
.map_err(|error| format!("serialize projected peer stream: {error}"))?;
if let Err(error) = projection_channel.send(NativeProjectionEvent {
request_id,
kind: "stream",
capability_keys: vec![capability_key],
payload,
}) {
state.peer_bi_streams.lock().await.remove(&result.stream_id);
state
.projected_stream_failures
.fetch_add(1, Ordering::Relaxed);
return Err(format!("send projected peer stream: {error}"));
}
state
.projected_stream_projected
.fetch_add(1, Ordering::Relaxed);
Ok(ProjectedPeerBiHandoff::Handled)
}
fn native_stream_trace_enabled() -> bool {
cfg!(debug_assertions)
|| std::env::var("OPENRTC_NATIVE_STREAM_TRACE").ok().as_deref() == Some("1")
}
fn spawn_peer_bi_stream_reader(
channel: Channel<Response>,
stream_id: String,
mut recv: PeerRecvStream,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut chunk = vec![0_u8; 64 * 1024];
let mut close_error: Option<String> = None;
loop {
match recv.read(&mut chunk).await {
Ok(0) => break,
Ok(n) => {
if native_stream_trace_enabled() {
eprintln!(
"[openrtc-tauri][peer-bi-stream] stream_id={} phase=plaintext-chunk bytes={}",
stream_id, n
);
}
if channel.send(Response::new(chunk[..n].to_vec())).is_err() {
close_error = Some("native stream IPC channel closed".to_string());
break;
}
}
Err(error) => {
if native_stream_trace_enabled() {
eprintln!(
"[openrtc-tauri][peer-bi-stream] stream_id={} phase=read-error error={}",
stream_id, error
);
}
close_error = Some(error.to_string());
break;
}
}
}
let close = serde_json::to_string(&PeerBiStreamClosedEvent {
r#type: "closed",
error: close_error,
})
.expect("peer stream close event serializes");
let _ = channel.send(Response::new(close));
})
}
async fn open_peer_bi_with<R: Runtime, F, Fut>(
_app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
open: F,
) -> Result<OpenBiResult, String>
where
F: FnOnce(Arc<openrtc::client::Client>, String, Option<u64>) -> Fut,
Fut: std::future::Future<
Output = anyhow::Result<(Option<String>, String, PeerSendStream, PeerRecvStream)>,
>,
{
let peer_id = peer_id.trim().to_string();
if peer_id.is_empty() {
return Err("peerId is required".to_string());
}
let (connection_id, remote_node_id, send, recv) = open(state.client(), peer_id, timeout_ms)
.await
.map_err(|error| format!("open peer bi stream failed: {error}"))?;
Ok(register_peer_bi_stream(&state, connection_id, remote_node_id, send, recv).await)
}
#[tauri::command]
async fn openrtc_set_identity_credential(
state: tauri::State<'_, OpenRtcTauriState>,
identity_credential: Option<String>,
) -> Result<(), String> {
state.set_identity_credential(identity_credential);
if *state.presence_loop_active.lock().await {
state.client().request_presence_update();
}
Ok(())
}
fn required_app_tag(value: &str) -> Result<&str, String> {
let value = value.trim();
if value.len() < 5
|| value.len() > 80
|| !value.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
})
{
return Err("OpenRTC app tag is invalid".to_string());
}
Ok(value)
}
fn validate_public_device_jwk(value: &serde_json::Value) -> Result<(), String> {
let object = value
.as_object()
.ok_or_else(|| "OpenRTC device public key is invalid".to_string())?;
let x = object
.get("x")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
if object.get("kty").and_then(serde_json::Value::as_str) != Some("OKP")
|| object.get("crv").and_then(serde_json::Value::as_str) != Some("Ed25519")
|| object.contains_key("d")
|| x.len() != 43
|| !x
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
{
return Err("OpenRTC device public key must be a public Ed25519 JWK".to_string());
}
Ok(())
}
fn required_secure_record_key(value: &str) -> Result<&str, String> {
let value = value.trim();
if value.is_empty()
|| value.len() > 512
|| !value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':'))
{
return Err("OpenRTC secure record key is invalid".to_string());
}
Ok(value)
}
#[tauri::command]
async fn openrtc_device_public_key(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
) -> Result<serde_json::Value, String> {
state.device_public_key(&app_tag)
}
#[tauri::command]
async fn openrtc_sign_device_proof(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
challenge: String,
) -> Result<String, String> {
state.sign_device_proof(&app_tag, &challenge)
}
#[tauri::command]
async fn openrtc_delete_device_key(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
) -> Result<(), String> {
state.delete_device_key(&app_tag)
}
#[tauri::command]
async fn openrtc_read_secure_record(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
key: String,
) -> Result<Option<String>, String> {
state.read_secure_record(&app_tag, &key)
}
#[tauri::command]
async fn openrtc_write_secure_record(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
key: String,
value: String,
) -> Result<(), String> {
state.write_secure_record(&app_tag, &key, &value)
}
#[tauri::command]
async fn openrtc_delete_secure_record(
state: tauri::State<'_, OpenRtcTauriState>,
app_tag: String,
key: String,
) -> Result<(), String> {
state.delete_secure_record(&app_tag, &key)
}
#[tauri::command]
async fn openrtc_offline_runtime_support(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<OfflineRuntimeSupport, String> {
Ok(state.offline_runtime_support())
}
#[tauri::command]
async fn openrtc_create_offline_enrollment_request(
state: tauri::State<'_, OpenRtcTauriState>,
trust_domain: String,
device_id: String,
endpoint_id: String,
enrollment_nonce: String,
requested_roles: Vec<String>,
) -> Result<openrtc::offline::OfflineEnrollmentRequest, String> {
let client = state.client();
let current_endpoint_id = client.current_node_id().await.ok_or_else(|| {
"OpenRTC Iroh endpoint must be started before offline enrollment".to_string()
})?;
if current_endpoint_id != endpoint_id.trim() {
return Err(
"offline enrollment endpoint does not match the native Rust endpoint".to_string(),
);
}
let app_tag = client.app_tag();
let signer = state
.native_device_key_signer
.as_deref()
.ok_or_else(|| "offline enrollment requires the native host device signer".to_string())?;
openrtc::offline::OfflineEnrollmentRequest::create(
&OfflineDeviceSigner { signer, app_tag },
&trust_domain,
&device_id,
¤t_endpoint_id,
&enrollment_nonce,
requested_roles,
signer.offline_assurance(app_tag),
openrtc::session_token::now_unix_ms(),
)
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn openrtc_verify_offline_enrollment_request(
request: openrtc::offline::OfflineEnrollmentRequest,
) -> Result<(), String> {
request
.verify()
.map(|_| ())
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn rtc_native_status(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<NativeRuntimeStatusResponse, String> {
Ok(state.project_public_runtime_status(state.client().runtime_status().await))
}
#[tauri::command]
async fn get_rtc_local_device_info<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
let data_dir = app_data_dir(&app, &state)?;
state
.client()
.init_native_device_identity(data_dir, None)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn update_rtc_local_device_name<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
device_name: String,
) -> Result<openrtc::native_device::NativeDeviceIdentity, String> {
let data_dir = app_data_dir(&app, &state)?;
let _ = state
.client()
.init_native_device_identity(data_dir, None)
.await
.map_err(|error| error.to_string())?;
state
.client()
.update_native_device_name(&device_name)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn get_iroh_node_id(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Option<String>, String> {
Ok(state.client().current_node_id().await)
}
#[tauri::command]
async fn start_iroh_node(state: tauri::State<'_, OpenRtcTauriState>) -> Result<String, String> {
ensure_iroh_node(state.client()).await
}
#[tauri::command]
async fn get_iroh_endpoint_ticket(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<String, String> {
ensure_iroh_node(state.client()).await?;
state
.client()
.endpoint_ticket()
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn register_session_token(
state: tauri::State<'_, OpenRtcTauriState>,
token: String,
scope: String,
max_connections: u32,
expires_at_ms: Option<u64>,
) -> Result<(), String> {
if let Some(expires_at_ms) = expires_at_ms {
state
.client()
.register_token_until(token, scope, max_connections, expires_at_ms);
} else {
state
.client()
.register_session_token(token, scope, max_connections);
}
Ok(())
}
#[tauri::command]
async fn get_endpoint_ticket_with_token(
state: tauri::State<'_, OpenRtcTauriState>,
scope: String,
max_connections: u32,
) -> Result<String, String> {
ensure_iroh_node(state.client()).await?;
state
.client()
.endpoint_ticket_with_token(&scope, max_connections)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn validate_session_token(
state: tauri::State<'_, OpenRtcTauriState>,
token: String,
connection_id: Option<String>,
) -> Result<String, String> {
if let Some(connection_id) = connection_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
state
.client()
.validate_connection_token(&token, connection_id, None)
.await
.map_err(|error| error.to_string())
} else {
state.client().validate_session_token(&token)
}
}
#[tauri::command]
async fn revoke_session_tokens_by_scope(
state: tauri::State<'_, OpenRtcTauriState>,
scope: String,
) -> Result<Vec<String>, String> {
let affected = state.client().revoke_tokens_by_scope(&scope).await;
if revokes_managed_session(&scope) {
if let Some(previous) = state.managed_session.lock().await.take() {
previous.owner_client.stop_presence_loop();
previous.owner_client.stop_auto_connect();
previous.owner_client.stop_external_auto_connect().await;
}
*state.presence_loop_active.lock().await = false;
}
Ok(affected)
}
#[tauri::command]
async fn stop_rtc_presence_loop(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
state.client().stop_presence_loop();
*state.presence_loop_active.lock().await = false;
if let Some(active) = state.managed_session.lock().await.as_mut() {
active.result.presence_started = false;
}
stop_subscription_by_id(&state, "internal-presence-loop").await;
Ok(())
}
#[tauri::command]
async fn register_rtc_capability(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
avenue_kind: String,
avenue_id: String,
sparse_fanout: Option<bool>,
architecture: Option<openrtc::native::RoomArchitectureMode>,
) -> Result<(), String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let avenue_kind = required_native_capability_part(avenue_kind, "avenue kind")?;
let avenue_id = required_native_capability_part(avenue_id, "avenue id")?;
if capability_key != format!("{avenue_kind}:{avenue_id}") {
return Err("native capability key does not match its avenue".to_string());
}
state
.capability_registry
.lock()
.await
.register_with_architecture(
capability_key,
avenue_kind,
avenue_id,
sparse_fanout.unwrap_or(false),
architecture,
)
}
#[cfg(feature = "managed-group-encryption")]
fn managed_group_state_path<R: Runtime>(
app: &tauri::AppHandle<R>,
state: &OpenRtcTauriState,
capability_key: &str,
device_id: &str,
) -> Result<PathBuf, String> {
let mut digest = Sha256::new();
digest.update(b"openrtc-tauri-managed-group-state-v1:");
digest.update(state.client().app_tag().as_bytes());
digest.update(device_id.as_bytes());
digest.update(capability_key.as_bytes());
let name = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(digest.finalize());
Ok(app_data_dir(app, state)?
.join("managed-room-state")
.join(format!("{name}.bin")))
}
#[cfg(feature = "managed-group-encryption")]
fn write_private_managed_group_state(path: &std::path::Path, bytes: &[u8]) -> Result<(), String> {
if bytes.is_empty() || bytes.len() > 16 * 1024 * 1024 {
return Err("managed room state exceeds its protected bound".to_string());
}
let parent = path
.parent()
.ok_or_else(|| "managed room state path has no parent".to_string())?;
std::fs::create_dir_all(parent)
.map_err(|error| format!("create managed room state directory: {error}"))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700))
.map_err(|error| format!("protect managed room state directory: {error}"))?;
}
let temporary = parent.join(format!(
".{}.{}.tmp",
path.file_name()
.and_then(|name| name.to_str())
.unwrap_or("state"),
uuid::Uuid::new_v4().simple(),
));
std::fs::write(&temporary, bytes)
.map_err(|error| format!("write managed room state: {error}"))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&temporary, std::fs::Permissions::from_mode(0o600))
.map_err(|error| format!("protect managed room state: {error}"))?;
}
std::fs::rename(&temporary, path)
.map_err(|error| format!("replace managed room state: {error}"))?;
Ok(())
}
#[cfg(feature = "managed-group-encryption")]
fn persist_managed_group_actions<R: Runtime>(
app: &tauri::AppHandle<R>,
state: &OpenRtcTauriState,
capability_key: &str,
device_id: &str,
actions: Vec<serde_json::Value>,
) -> Result<Vec<serde_json::Value>, String> {
let mut outbound = Vec::with_capacity(actions.len());
for action in actions {
if action.get("type").and_then(serde_json::Value::as_str) == Some("persist-state") {
let sealed = action
.get("sealedState")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "managed room persisted action is invalid".to_string())?;
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(sealed)
.map_err(|_| "managed room persisted action encoding is invalid".to_string())?;
write_private_managed_group_state(
&managed_group_state_path(app, state, capability_key, device_id)?,
&bytes,
)?;
} else {
outbound.push(action);
}
}
Ok(outbound)
}
#[cfg(feature = "managed-group-encryption")]
fn derive_tauri_managed_group_key(
state: &OpenRtcTauriState,
capability_key: &str,
device_id: &str,
) -> Result<[u8; 32], String> {
let client = state.client();
let app_tag = client.app_tag();
let challenge =
format!("openrtc:managed-group-wrapping-key:v1:{app_tag}:{device_id}:{capability_key}",);
let signer = state.native_device_key_signer.as_ref().ok_or_else(|| {
"managed rooms require the host's sign-only native device key".to_string()
})?;
let mut signature = signer.sign(app_tag, challenge.as_bytes())?;
if signature.len() != 64 {
return Err("native device signer returned an invalid managed-room signature".to_string());
}
let mut digest = Sha256::new();
digest.update(b"openrtc:managed-group-key-derivation:v1:");
digest.update(challenge.as_bytes());
digest.update(&signature);
signature.fill(0);
Ok(digest.finalize().into())
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
async fn openrtc_init_managed_room_group<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
device_id: String,
) -> Result<serde_json::Value, String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let device_id = required_native_capability_part(device_id, "device id")?;
{
let registry = state.capability_registry.lock().await;
let registration = registry
.registrations
.get(&capability_key)
.ok_or_else(|| "managed room capability is not registered".to_string())?;
if registration.avenue_kind != "room" {
return Err("managed fanout is available only for room capabilities".to_string());
}
}
let path = managed_group_state_path(&app, &state, &capability_key, &device_id)?;
let sealed = match std::fs::read(path) {
Ok(bytes) if bytes.len() <= 16 * 1024 * 1024 => Some(bytes),
Ok(_) => return Err("managed room persisted state exceeds its protected bound".to_string()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
Err(error) => return Err(format!("read managed room persisted state: {error}")),
};
let controller = openrtc::native::NativeManagedGroupController::new(
&device_id,
derive_tauri_managed_group_key(&state, &capability_key, &device_id)?,
sealed.as_deref(),
)
.map_err(|error| format!("initialize managed room group: {error:#}"))?;
let action = controller
.publish_key_package()
.map_err(|error| format!("create managed room KeyPackage: {error:#}"))?;
state
.managed_groups
.lock()
.await
.insert(capability_key, controller);
Ok(action)
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
async fn openrtc_handle_managed_room_prepare_page<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
device_id: String,
page: serde_json::Value,
) -> Result<Vec<serde_json::Value>, String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let device_id = required_native_capability_part(device_id, "device id")?;
let actions = state
.managed_groups
.lock()
.await
.get_mut(&capability_key)
.ok_or_else(|| "managed room group is not initialized".to_string())?
.handle_prepare_page(page)
.map_err(|error| format!("handle managed room preparation: {error:#}"))?;
persist_managed_group_actions(&app, &state, &capability_key, &device_id, actions)
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
async fn openrtc_handle_managed_room_artifact_chunk<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
device_id: String,
chunk: serde_json::Value,
) -> Result<Vec<serde_json::Value>, String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let device_id = required_native_capability_part(device_id, "device id")?;
let actions = state
.managed_groups
.lock()
.await
.get_mut(&capability_key)
.ok_or_else(|| "managed room group is not initialized".to_string())?
.handle_artifact_chunk(chunk)
.map_err(|error| format!("handle managed room artifact: {error:#}"))?;
persist_managed_group_actions(&app, &state, &capability_key, &device_id, actions)
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn openrtc_seal_managed_room_payload<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
device_id: String,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: String,
channel: String,
priority: u8,
zone_id: Option<String>,
payload: Vec<u8>,
) -> Result<openrtc::native::NativeManagedProtectedPayload, String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let device_id = required_native_capability_part(device_id, "device id")?;
let protected = state
.managed_groups
.lock()
.await
.get_mut(&capability_key)
.ok_or_else(|| "managed room group is not initialized".to_string())?
.seal_payload(
architecture_epoch,
encryption_epoch,
&message_id,
&channel,
priority,
zone_id.as_deref(),
&payload,
)
.map_err(|error| format!("seal managed room payload: {error:#}"))?;
let sealed = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(&protected.sealed_state)
.map_err(|_| "managed room state encoding is invalid".to_string())?;
write_private_managed_group_state(
&managed_group_state_path(&app, &state, &capability_key, &device_id)?,
&sealed,
)?;
Ok(protected)
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn openrtc_open_managed_room_payload<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
device_id: String,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: String,
channel: String,
priority: u8,
zone_id: Option<String>,
ciphertext: Vec<u8>,
) -> Result<openrtc::native::NativeManagedProtectedPayload, String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let device_id = required_native_capability_part(device_id, "device id")?;
let protected = state
.managed_groups
.lock()
.await
.get_mut(&capability_key)
.ok_or_else(|| "managed room group is not initialized".to_string())?
.open_payload(
architecture_epoch,
encryption_epoch,
&message_id,
&channel,
priority,
zone_id.as_deref(),
&ciphertext,
)
.map_err(|error| format!("open managed room payload: {error:#}"))?;
let sealed = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(&protected.sealed_state)
.map_err(|_| "managed room state encoding is invalid".to_string())?;
write_private_managed_group_state(
&managed_group_state_path(&app, &state, &capability_key, &device_id)?,
&sealed,
)?;
Ok(protected)
}
#[cfg(feature = "managed-group-encryption")]
#[tauri::command]
async fn openrtc_forget_managed_room_group(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
) -> Result<(), String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
state.managed_groups.lock().await.remove(&capability_key);
Ok(())
}
#[cfg(not(feature = "managed-group-encryption"))]
macro_rules! unavailable_managed_room_command {
($name:ident ( $($arg:ident : $type:ty),* ) -> $return:ty) => {
#[tauri::command]
async fn $name($($arg: $type),*) -> Result<$return, String> {
$(let _ = $arg;)*
Err("native managed rooms were not compiled into this Tauri host".to_string())
}
};
}
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_init_managed_room_group(
capability_key: String, device_id: String
) -> serde_json::Value);
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_handle_managed_room_prepare_page(
capability_key: String, device_id: String, page: serde_json::Value
) -> Vec<serde_json::Value>);
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_handle_managed_room_artifact_chunk(
capability_key: String, device_id: String, chunk: serde_json::Value
) -> Vec<serde_json::Value>);
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_seal_managed_room_payload(
capability_key: String, device_id: String, architecture_epoch: u64,
encryption_epoch: u64, message_id: String, channel: String, priority: u8,
zone_id: Option<String>, payload: Vec<u8>
) -> serde_json::Value);
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_open_managed_room_payload(
capability_key: String, device_id: String, architecture_epoch: u64,
encryption_epoch: u64, message_id: String, channel: String, priority: u8,
zone_id: Option<String>, ciphertext: Vec<u8>
) -> serde_json::Value);
#[cfg(not(feature = "managed-group-encryption"))]
unavailable_managed_room_command!(openrtc_forget_managed_room_group(
capability_key: String
) -> ());
#[tauri::command]
async fn unregister_rtc_capability(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
) -> Result<(), String> {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let (empty, aggregate) = {
let mut registry = state.capability_registry.lock().await;
registry.registrations.remove(&capability_key);
let empty = registry.registrations.is_empty();
let aggregate = (!empty)
.then(|| registry.aggregate_desired_peers())
.transpose()?;
(empty, aggregate)
};
if empty {
state.client().stop_external_auto_connect().await;
} else if let Some((revision, peers_json)) = aggregate {
state
.client()
.submit_external_desired_peers(revision, &peers_json)
.await
.map_err(|error| error.to_string())?;
}
Ok(())
}
#[tauri::command]
async fn start_rtc_managed_session<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
device_name: Option<String>,
local_device_id: Option<String>,
metadata: Option<String>,
transports: Option<openrtc::client::TransportConfig>,
auto_connect: Option<bool>,
presence: Option<bool>,
) -> Result<StartSessionResult, String> {
let command_started = std::time::Instant::now();
let _start_guard = state.managed_session_start_guard.lock().await;
let client = state.client();
let app_tag = client.app_tag().to_string();
eprintln!(
"[openrtc-tauri][managed-session] start user_id={} app_tag={} has_device_name={} has_local_device_id={} has_metadata={} auto_connect={} presence={}",
user_id,
app_tag,
device_name
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
local_device_id
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
metadata
.as_deref()
.map(str::trim)
.map(|value| !value.is_empty())
.unwrap_or(false),
auto_connect.unwrap_or(false),
presence.unwrap_or(false)
);
let data_dir = app_data_dir(&app, &state)?;
let install_context = InstallContext {
data_dir: data_dir.clone(),
};
ensure_requested_native_transports(&state, &client, transports.as_ref(), &install_context)
.await?;
if let Some(transport_config) = transports {
client
.update_transport_config(transport_config)
.await
.map_err(|error| error.to_string())?;
}
let local_node_id = ensure_iroh_node(client.clone()).await?;
let local_device = client
.init_native_device_identity(data_dir, device_name.as_deref())
.await
.map_err(|error| error.to_string())?;
let effective_local_device_id =
managed_session_device_id(local_device_id.as_deref(), &local_device);
if let Some(requested) = requested_local_device_id(local_device_id.as_deref()) {
if requested != local_device.device_id {
eprintln!(
"[openrtc-tauri] requested localDeviceId={} is an alias; persisted native identity {} remains authoritative",
requested, local_device.device_id
);
}
}
let should_start_presence = presence.unwrap_or(false);
let should_start_auto_connect = auto_connect.unwrap_or(false);
if should_start_presence || should_start_auto_connect {
return Err(
"native Firebase coordination has been removed; publish gateway presence and submit external desired peers instead"
.to_string(),
);
}
let managed_session_key = format!("{}:{}", app_tag, effective_local_device_id);
let (disposition, reused) = {
let active = state.managed_session.lock().await;
let disposition = managed_session_disposition(
active.as_ref().map(|record| record.key.as_str()),
active
.as_ref()
.is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
active
.as_ref()
.is_some_and(|record| record.result.presence_started),
active
.as_ref()
.is_some_and(|record| record.result.auto_connect_started),
&managed_session_key,
should_start_presence,
should_start_auto_connect,
);
let reused = (disposition == SessionDisposition::Reuse).then(|| {
active
.as_ref()
.expect("managed session reuse requires an active record")
.result
.clone()
});
(disposition, reused)
};
if let Some(reused) = reused {
eprintln!(
"[openrtc-tauri][managed-session] reused key={} elapsed_ms={}",
managed_session_key,
command_started.elapsed().as_millis()
);
return Ok(reused);
}
if disposition == SessionDisposition::Replace {
let previous = state
.managed_session
.lock()
.await
.take()
.expect("managed session replacement requires an active record");
previous.owner_client.stop_presence_loop();
previous.owner_client.stop_auto_connect();
previous.owner_client.stop_external_auto_connect().await;
*state.presence_loop_active.lock().await = false;
}
let ticket_scope = Some("user-device".to_string());
let ticket = client
.endpoint_ticket_with_token("user-device", 0)
.await
.map_err(|error| error.to_string())?;
let logical_local_device = local_device.clone();
let owner_epoch = match disposition {
SessionDisposition::Refresh => {
state
.managed_session
.lock()
.await
.as_ref()
.expect("managed session refresh requires an active record")
.owner_epoch
}
SessionDisposition::Start | SessionDisposition::Replace => {
state.allocate_managed_session_owner_epoch()
}
SessionDisposition::Reuse => unreachable!("reuse returns before session startup"),
};
eprintln!(
"[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={}",
user_id,
app_tag,
local_node_id,
logical_local_device.device_id,
owner_epoch,
should_start_presence,
should_start_auto_connect,
ticket_scope.as_deref().unwrap_or("unrestricted"),
command_started.elapsed().as_millis()
);
let result = StartSessionResult {
local_node_id,
ticket_scope,
ticket: Some(ticket),
presence_started: false,
auto_connect_started: false,
local_device: logical_local_device,
};
*state.managed_session.lock().await = Some(ManagedSessionRecord {
key: managed_session_key,
owner_client: client,
owner_epoch,
result: result.clone(),
});
Ok(result)
}
#[tauri::command]
async fn start_rtc_external_auto_connect(
state: tauri::State<'_, OpenRtcTauriState>,
user_id: String,
local_device_id: String,
capability_key: Option<String>,
) -> Result<(), String> {
if let Some(capability_key) = capability_key {
let capability_key = required_native_capability_part(capability_key, "capability key")?;
if !state
.capability_registry
.lock()
.await
.registrations
.contains_key(&capability_key)
{
return Err(format!(
"native capability {capability_key} is not registered"
));
}
}
let root_user_id = format!("{}:native-root", state.client().app_tag());
let _requested_capability_principal = user_id;
state
.client()
.start_external_auto_connect(root_user_id, local_device_id)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn submit_rtc_desired_peers(
state: tauri::State<'_, OpenRtcTauriState>,
revision: u64,
peers_json: String,
capability_key: Option<String>,
local_device_id: Option<String>,
app_tag: Option<String>,
) -> Result<bool, String> {
let original_peers = serde_json::from_str::<Vec<serde_json::Value>>(&peers_json)
.map_err(|error| format!("invalid capability desired-peer payload: {error}"))?;
if original_peers.len() > 100 {
return Err("capability desired-peer payload exceeds 100 peers".to_string());
}
let key = {
let registry = state.capability_registry.lock().await;
match capability_key {
Some(key) => required_native_capability_part(key, "capability key")?,
None if registry.registrations.len() == 1 => registry
.registrations
.keys()
.next()
.cloned()
.expect("one registration has one key"),
None => {
return Err(
"capabilityKey is required when multiple native capabilities are active"
.to_string(),
)
}
}
};
let sparse = state
.capability_registry
.lock()
.await
.registrations
.get(&key)
.ok_or_else(|| format!("native capability {key} is not registered"))?
.uses_sparse_fanout();
let peers = if sparse {
let local_device_id = local_device_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| "sparse native capability requires localDeviceId".to_string())?;
let app_tag = app_tag
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| "sparse native capability requires appTag".to_string())?;
let local_public_jwk = state.device_public_key(app_tag)?;
let local_device_key_x = local_public_jwk
.get("x")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "native sparse fanout signer returned no public key".to_string())?;
let projected = state
.client()
.configure_sparse_fanout(
&key,
revision,
local_device_id,
local_device_key_x,
&peers_json,
true,
)
.await
.map_err(|error| error.to_string())?
.0;
serde_json::from_str(&projected)
.map_err(|error| format!("invalid projected desired-peer payload: {error}"))?
} else {
original_peers
};
let (root_revision, aggregate) = {
let mut registry = state.capability_registry.lock().await;
let registration = registry
.registrations
.get_mut(&key)
.ok_or_else(|| format!("native capability {key} is not registered"))?;
if revision <= registration.desired_revision {
return Ok(false);
}
registration.desired_revision = revision;
registration.desired_peers = peers;
registry.aggregate_desired_peers()?
};
let accepted = state
.client()
.submit_external_desired_peers(root_revision, &aggregate)
.await
.map_err(|error| error.to_string())?;
if accepted {
replay_current_connection_states(&state).await;
}
Ok(accepted)
}
#[tauri::command]
async fn stop_rtc_auto_connect(state: tauri::State<'_, OpenRtcTauriState>) -> Result<(), String> {
state.client().stop_auto_connect();
state.client().stop_external_auto_connect().await;
if let Some(active) = state.managed_session.lock().await.as_mut() {
active.result.auto_connect_started = false;
}
Ok(())
}
#[tauri::command]
async fn connect_to_device(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: Option<String>,
endpoint_ticket: String,
timeout_ms: Option<u64>,
) -> Result<openrtc::client::ManagedConnectResult, String> {
let client = state.client();
let connect = client.connect_device(device_id.as_deref(), &endpoint_ticket);
match timeout_ms {
Some(timeout_ms) => {
tokio::time::timeout(std::time::Duration::from_millis(timeout_ms.max(1)), connect)
.await
.map_err(|_| format!("connect_to_device timeout after {timeout_ms}ms"))?
.map_err(|error| error.to_string())
}
None => connect.await.map_err(|error| error.to_string()),
}
}
#[tauri::command]
async fn disconnect_device(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
node_id_hint: Option<String>,
) -> Result<(), String> {
state
.client()
.disconnect_device(&device_id, node_id_hint.as_deref())
.await;
Ok(())
}
#[tauri::command]
async fn set_auto_connect_excluded(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
excluded: bool,
) -> Result<(), String> {
if excluded {
state.client().exclude_peer_and_publish(&device_id).await;
} else {
state.client().unexclude_peer_and_publish(&device_id).await;
}
Ok(())
}
#[tauri::command]
async fn set_rtc_external_auto_connect_excluded(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
excluded: bool,
) -> Result<(), String> {
state
.client()
.set_auto_connect_excluded(&device_id, excluded);
state.client().wake_native_external_auto_connect().await;
Ok(())
}
#[tauri::command]
async fn resolve_rtc_peer_connection_records(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
Ok(state.client().resolve_peer_connection_records(&id).await)
}
#[tauri::command]
async fn resolve_rtc_peer_identity(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Option<openrtc::connection_manager::PeerSnapshot>, String> {
Ok(state.client().peer_snapshot(&id).await)
}
#[tauri::command]
async fn get_rtc_peer_session(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().peer_session(&id).await)
}
#[tauri::command]
async fn list_rtc_peer_sessions(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().peer_sessions().await)
}
#[tauri::command]
async fn list_rtc_managed_connections(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<Vec<openrtc::connection_manager::ConnectionRecord>, String> {
Ok(state.client().list_managed_connections().await)
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct NotifyRtcNetworkChangeResult {
retired_stale_connections: usize,
}
#[tauri::command]
async fn notify_rtc_network_change(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<NotifyRtcNetworkChangeResult, String> {
let retired_stale_connections = state
.client()
.notify_network_change()
.await
.map_err(|error| error.to_string())?;
Ok(NotifyRtcNetworkChangeResult {
retired_stale_connections,
})
}
#[tauri::command]
async fn set_rtc_transport_priority(
state: tauri::State<'_, OpenRtcTauriState>,
priority: Vec<openrtc::route_policy::KnownRoute>,
) -> Result<(), String> {
let client = state.client();
let mut transport_config = client.transport_config().await;
transport_config.route_priority = priority;
client
.update_transport_config(transport_config)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn wait_for_rtc_settled_peer(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
timeout_ms: Option<u64>,
) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().wait_for_peer(&id, timeout_ms).await)
}
#[tauri::command]
async fn list_rtc_connection_states(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
) -> Result<Vec<openrtc::client::StateSnapshot>, String> {
let states = state.client().connection_states().await;
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let registry = state.capability_registry.lock().await;
Ok(states
.into_iter()
.filter(|snapshot| {
registry
.capability_keys_for_state(snapshot)
.iter()
.any(|key| key == &capability_key)
})
.collect())
}
#[tauri::command]
async fn get_rtc_connection_state(
state: tauri::State<'_, OpenRtcTauriState>,
connection_id: String,
capability_key: String,
) -> Result<Option<openrtc::client::StateSnapshot>, String> {
let snapshot = state.client().connection_state(&connection_id).await;
let Some(snapshot) = snapshot else {
return Ok(None);
};
let capability_key = required_native_capability_part(capability_key, "capability key")?;
let included = state
.capability_registry
.lock()
.await
.capability_keys_for_state(&snapshot)
.iter()
.any(|key| key == &capability_key);
Ok(included.then_some(snapshot))
}
#[tauri::command]
async fn stop_rtc_subscription(
state: tauri::State<'_, OpenRtcTauriState>,
request_id: String,
) -> Result<(), String> {
stop_subscription_by_id(&state, &request_id).await;
Ok(())
}
async fn stop_subscription_by_id(state: &OpenRtcTauriState, request_id: &str) {
if let Some(handle) = state.subscriptions.lock().await.remove(request_id) {
handle.abort();
}
let mut projected = state.native_projection_subscription.lock().await;
if projected
.as_ref()
.is_some_and(|subscription| subscription.request_id == request_id)
{
*projected = None;
}
}
#[tauri::command]
async fn start_native_projection(
state: tauri::State<'_, OpenRtcTauriState>,
request_id: Option<String>,
channel: Channel<NativeProjectionEvent>,
) -> Result<String, String> {
let request_id = request_id
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
let mut subscription = state.native_projection_subscription.lock().await;
if let Some(current) = subscription.as_ref() {
if current.request_id == request_id {
return Ok(request_id);
}
return Err(
"OpenRTC native projection subscription already has a process owner".to_string(),
);
}
*subscription = Some(NativeProjectionSubscription {
request_id: request_id.clone(),
channel,
});
Ok(request_id)
}
#[tauri::command]
async fn get_projected_stream_diagnostics(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<ProjectedStreamDiagnostics, String> {
Ok(ProjectedStreamDiagnostics {
subscription_active: state.native_projection_subscription.lock().await.is_some(),
offered: state.projected_stream_offered.load(Ordering::Relaxed),
decoded: state.projected_stream_decoded.load(Ordering::Relaxed),
projected: state.projected_stream_projected.load(Ordering::Relaxed),
unhandled_non_channel: state
.projected_stream_unhandled_non_channel
.load(Ordering::Relaxed),
unhandled_other_channel: state
.projected_stream_unhandled_other_channel
.load(Ordering::Relaxed),
unhandled_no_subscription: state
.projected_stream_unhandled_no_subscription
.load(Ordering::Relaxed),
failures: state.projected_stream_failures.load(Ordering::Relaxed),
last_transport_stable_id: state
.projected_stream_last_transport_stable_id
.load(Ordering::Relaxed),
validation_checks: state
.projected_stream_validation_checks
.load(Ordering::Relaxed),
validation_rejections: state
.projected_stream_validation_rejections
.load(Ordering::Relaxed),
last_validated_transport_stable_id: state
.projected_stream_last_validated_transport_stable_id
.load(Ordering::Relaxed),
authorized: state.projected_stream_authorized.load(Ordering::Relaxed),
unauthorized: state.projected_stream_unauthorized.load(Ordering::Relaxed),
last_channel: state.projected_stream_last_channel.lock().await.clone(),
last_protocol: state.projected_stream_last_protocol.lock().await.clone(),
last_connection_id: state
.projected_stream_last_connection_id
.lock()
.await
.clone(),
last_remote_node_id: state
.projected_stream_last_remote_node_id
.lock()
.await
.clone(),
})
}
#[tauri::command]
async fn is_current_transport_stable_id(
state: tauri::State<'_, OpenRtcTauriState>,
endpoint_id: String,
transport_stable_id: u64,
) -> Result<bool, String> {
state
.projected_stream_validation_checks
.fetch_add(1, Ordering::Relaxed);
state
.projected_stream_last_validated_transport_stable_id
.store(transport_stable_id, Ordering::Relaxed);
let is_current = state
.client()
.is_current_transport_stable_id_str(&endpoint_id, transport_stable_id)
.await
.map_err(|error| error.to_string())?;
if !is_current {
state
.projected_stream_validation_rejections
.fetch_add(1, Ordering::Relaxed);
}
if native_stream_trace_enabled() {
eprintln!(
"[openrtc-tauri][incoming-application-stream] remote_node_id={} transport_stable_id={} current={} phase=validate",
endpoint_id, transport_stable_id, is_current
);
}
Ok(is_current)
}
#[tauri::command]
async fn send_peer_message(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<(), String> {
let id = required_request_header(&request, PEER_ID_HEADER)?;
let data = request_body_bytes(&request)?;
state
.client()
.send_peer(&id, &data)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn encode_sparse_fanout_message(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<Response, String> {
let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
let app_tag = required_request_header(&request, FANOUT_APP_TAG_HEADER)?;
let payload = request_body_bytes(&request)?;
let signing_request = state
.client()
.prepare_sparse_fanout_message(&capability, &payload)
.await
.map_err(|error| error.to_string())?;
let signature = state.sign_device_message(&app_tag, &signing_request.signing_input)?;
state
.client()
.finalize_sparse_fanout_message(&signing_request.request_id, &signature)
.await
.map(Response::new)
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn accept_sparse_fanout_message(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<serde_json::Value, String> {
let capability = required_request_header(&request, FANOUT_CAPABILITY_HEADER)?;
let source_peer = required_request_header(&request, FANOUT_SOURCE_PEER_HEADER)?;
let encoded = request_body_bytes(&request)?;
serde_json::to_value(
state
.client()
.accept_sparse_fanout_message(&capability, &source_peer, &encoded)
.await
.map_err(|error| error.to_string())?,
)
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn sparse_fanout_diagnostics(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
) -> Result<openrtc::sparse_fanout::SparseFanoutDiagnostics, String> {
Ok(state
.client()
.sparse_fanout_diagnostics(&required_native_capability_part(
capability_key,
"capability key",
)?)
.await)
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeBroadcastPublisherDescriptor {
handle: String,
verification_key: Vec<u8>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeBroadcastSessionDescriptor {
handle_id: String,
id: String,
role: openrtc::broadcast::BroadcastRole,
state: openrtc::broadcast::BroadcastState,
budget: openrtc::broadcast::BroadcastBudget,
}
#[tauri::command]
async fn prepare_openrtc_broadcast_publisher(
state: tauri::State<'_, OpenRtcTauriState>,
) -> Result<NativeBroadcastPublisherDescriptor, String> {
let signer = openrtc::broadcast::BroadcastPublisherSigner::generate()
.map_err(|error| error.to_string())?;
let verification_key = signer.verifying_key().to_vec();
let handle = uuid::Uuid::new_v4().simple().to_string();
state
.broadcast_signers
.lock()
.await
.insert(handle.clone(), signer);
Ok(NativeBroadcastPublisherDescriptor {
handle,
verification_key,
})
}
#[tauri::command]
async fn release_openrtc_broadcast_publisher(
state: tauri::State<'_, OpenRtcTauriState>,
handle: String,
) -> Result<bool, String> {
Ok(state
.broadcast_signers
.lock()
.await
.remove(&handle)
.is_some())
}
async fn drive_native_broadcast_actions(
session: &openrtc::broadcast::BroadcastSession,
adapter: &dyn NativeBroadcastAdapter,
) -> Result<(), String> {
loop {
let actions = session.take_actions(openrtc::broadcast::MAX_BROADCAST_ACTIONS);
if actions.is_empty() {
return Ok(());
}
for action in actions {
if let Some(observation) = adapter.apply(session.clone(), action).await? {
session
.observe(observation)
.map_err(|error| error.to_string())?;
}
}
}
}
async fn native_broadcast_driver(
state: &OpenRtcTauriState,
handle_id: &str,
) -> Result<Arc<Mutex<()>>, String> {
state
.broadcast_drivers
.lock()
.await
.get(handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())
}
#[tauri::command]
async fn open_openrtc_broadcast(
state: tauri::State<'_, OpenRtcTauriState>,
grant_token: String,
issuer_public_key: Vec<u8>,
publisher_signer_handle: Option<String>,
now_ms: u64,
) -> Result<NativeBroadcastSessionDescriptor, String> {
let issuer_bytes: [u8; 32] = issuer_public_key
.try_into()
.map_err(|_| "broadcast issuer key must be 32 bytes".to_string())?;
let issuer = VerifyingKey::from_bytes(&issuer_bytes)
.map_err(|_| "broadcast issuer key is invalid".to_string())?;
#[cfg(feature = "native-broadcast-moq")]
let publication_verifying_key = if let Some(handle) = publisher_signer_handle.as_deref() {
Some(
state
.broadcast_signers
.lock()
.await
.get(handle)
.ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?
.verifying_key(),
)
} else {
None
};
let broadcasts = state.client().broadcasts();
let challenge = broadcasts
.prepare_grant_verification(&grant_token, &issuer, now_ms)
.map_err(|error| error.to_string())?;
let device_signer = state.native_device_key_signer.as_deref().ok_or_else(|| {
"native managed broadcast requires the host device-key signer".to_string()
})?;
let binding_signature = device_signer
.sign(state.client().app_tag(), &challenge.signing_bytes)
.map_err(|error| format!("sign broadcast installation challenge: {error}"))?;
let grant = broadcasts
.complete_grant_verification(
&grant_token,
&issuer,
&challenge.handle,
&binding_signature,
now_ms,
)
.map_err(|error| error.to_string())?;
let grant_generation = grant.grant_generation();
#[cfg(feature = "native-broadcast-moq")]
let native_moq_access = if state.native_broadcast_moq_adapter.is_some() {
Some(
state
.resolve_native_broadcast_access(
&grant_token,
publication_verifying_key.as_ref(),
&issuer_bytes,
)
.await?,
)
} else {
None
};
let session = if let Some(handle) = publisher_signer_handle.as_deref() {
let signers = state.broadcast_signers.lock().await;
let signer = signers
.get(handle)
.ok_or_else(|| "broadcast publisher signer handle is invalid".to_string())?;
state
.client()
.broadcasts()
.open_publisher(grant, signer, now_ms)
} else {
state.client().broadcasts().open(grant, now_ms)
}
.map_err(|error| error.to_string())?;
let handle_id = format!("{}:{grant_generation}", session.id());
#[cfg(feature = "native-broadcast-moq")]
if let (Some(adapter), Some(access)) = (
state.native_broadcast_moq_adapter.as_ref(),
native_moq_access,
) {
adapter
.authorize(session.id(), grant_generation, access)
.await;
}
let adapter = state.native_broadcast_adapter.as_deref().ok_or_else(|| {
session.close();
"native managed broadcast adapter is unavailable".to_string()
})?;
if let Err(error) = drive_native_broadcast_actions(&session, adapter).await {
session.close();
#[cfg(feature = "native-broadcast-moq")]
if let Some(adapter) = state.native_broadcast_moq_adapter.as_ref() {
adapter
.forget_authorization(&session.id(), grant_generation)
.await;
}
return Err(error);
}
let descriptor = NativeBroadcastSessionDescriptor {
handle_id: handle_id.clone(),
id: session.id(),
role: session.role(),
state: session.state(),
budget: session.budget(),
};
state
.broadcast_sessions
.lock()
.await
.insert(handle_id.clone(), session);
state
.broadcast_drivers
.lock()
.await
.insert(handle_id, Arc::new(Mutex::new(())));
Ok(descriptor)
}
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn begin_openrtc_broadcast_publication(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
source_slot: String,
publication_id: Option<String>,
kind: String,
codec: String,
clock_rate: u32,
coded_width: Option<u32>,
coded_height: Option<u32>,
channels: Option<u16>,
) -> Result<PortableMediaPublicationStart, String> {
let driver = native_broadcast_driver(&state, &handle_id).await?;
let _driver = driver.lock().await;
let publication_id = match publication_id.as_deref().map(str::trim) {
Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
_ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
};
let kind: openrtc::media::MediaKind = kind
.parse()
.map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
if !matches!(
kind,
openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
) {
return Err("broadcast publications must be audio or video".to_string());
}
let publication = openrtc::media::MediaPublicationConfig {
publication_id,
media_generation: 0,
kind,
codec: codec
.parse()
.map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?,
clock_rate,
coded_width,
coded_height,
channels,
};
let session = state
.broadcast_sessions
.lock()
.await
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
let publication = session
.begin_publication(&source_slot, publication)
.map_err(|error| error.to_string())?;
let adapter = state
.native_broadcast_adapter
.as_deref()
.ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
drive_native_broadcast_actions(&session, adapter).await?;
Ok(PortableMediaPublicationStart {
publication_id: publication.publication_id.to_string(),
media_generation: publication.media_generation,
control: Vec::new(),
})
}
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn publish_openrtc_broadcast_sample(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
publication_id: String,
timestamp_us: u64,
duration_us: u32,
keyframe: bool,
discardable: bool,
payload: Vec<u8>,
) -> Result<(), String> {
let driver = native_broadcast_driver(&state, &handle_id).await?;
let _driver = driver.lock().await;
openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
.map_err(|error| error.to_string())?;
let session = state
.broadcast_sessions
.lock()
.await
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
session
.publish(
parse_media_publication_id(&publication_id)?,
openrtc::media::EncodedMediaSample {
timestamp_us,
duration_us,
keyframe,
discardable,
payload,
},
)
.map_err(|error| error.to_string())?;
let adapter = state
.native_broadcast_adapter
.as_deref()
.ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
drive_native_broadcast_actions(&session, adapter).await
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeBroadcastMediaSample {
publication_id: String,
media_generation: u32,
sequence: u64,
timestamp_us: u64,
duration_us: u32,
kind: &'static str,
codec: &'static str,
keyframe: bool,
discardable: bool,
payload: Vec<u8>,
}
struct NativeBroadcastReceivedObject {
source_slot: Option<String>,
object: Vec<u8>,
}
#[tauri::command]
async fn receive_openrtc_broadcast_media(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
max: usize,
) -> Result<Vec<NativeBroadcastMediaSample>, String> {
let driver = native_broadcast_driver(&state, &handle_id).await?;
let _driver = driver.lock().await;
let session = state
.broadcast_sessions
.lock()
.await
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
let adapter = state
.native_broadcast_adapter
.as_deref()
.ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
drive_native_broadcast_actions(&session, adapter).await?;
let objects: Vec<NativeBroadcastReceivedObject>;
#[cfg(feature = "native-broadcast-moq")]
{
objects = if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
native_moq
.take_inbound_objects_with_source(&session, max.min(64))
.await?
.into_iter()
.map(|item| NativeBroadcastReceivedObject {
source_slot: Some(item.source_slot),
object: item.object,
})
.collect()
} else {
adapter
.take_inbound_objects(session.clone(), max.min(64))
.await?
.into_iter()
.map(|object| NativeBroadcastReceivedObject {
source_slot: None,
object,
})
.collect()
};
}
#[cfg(not(feature = "native-broadcast-moq"))]
{
objects = adapter
.take_inbound_objects(session.clone(), max.min(64))
.await?
.into_iter()
.map(|object| NativeBroadcastReceivedObject {
source_slot: None,
object,
})
.collect();
}
let mut accepted = Vec::with_capacity(objects.len());
let mut delivered_bytes = 0_u64;
for received in objects {
let object_len = received.object.len() as u64;
let chunk = if let Some(source_slot) = received.source_slot.as_deref() {
let Some(media) = session
.accept_media_for_source(source_slot, &received.object)
.map_err(|error| error.to_string())?
else {
continue;
};
delivered_bytes = delivered_bytes.saturating_add(object_len);
match media.media {
Some(media) => Some(
openrtc::media::EncodedMediaChunk::decode(&media)
.map_err(|error| error.to_string())?,
),
None => None,
}
} else {
let chunk = session
.accept_object(&received.object)
.map_err(|error| error.to_string())?;
delivered_bytes = delivered_bytes.saturating_add(object_len);
chunk
};
let Some(chunk) = chunk else { continue };
openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
.and_then(|_| {
openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
})
.map_err(|error| error.to_string())?;
accepted.push(NativeBroadcastMediaSample {
publication_id: chunk.publication_id.to_string(),
media_generation: chunk.media_generation,
sequence: chunk.sequence,
timestamp_us: chunk.timestamp_us,
duration_us: chunk.duration_us,
kind: media_kind_name(chunk.kind),
codec: media_codec_name(chunk.codec),
keyframe: chunk.keyframe,
discardable: chunk.discardable,
payload: chunk.payload,
});
}
#[cfg(feature = "native-broadcast-moq")]
if let Some(native_moq) = state.native_broadcast_moq_adapter.as_ref() {
let usage = native_moq
.record_delivered_bytes(&session, delivered_bytes)
.await;
let settlement = drive_native_broadcast_actions(&session, adapter).await;
usage?;
settlement?;
}
Ok(accepted)
}
#[tauri::command]
async fn revoke_openrtc_broadcast(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
grant_generation: u64,
) -> Result<(), String> {
let driver = native_broadcast_driver(&state, &handle_id).await?;
let _driver = driver.lock().await;
let session = state
.broadcast_sessions
.lock()
.await
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
session
.revoke(grant_generation)
.map_err(|error| error.to_string())?;
let adapter = state
.native_broadcast_adapter
.as_deref()
.ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
drive_native_broadcast_actions(&session, adapter).await
}
#[tauri::command]
async fn close_openrtc_broadcast(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
) -> Result<(), String> {
let driver = native_broadcast_driver(&state, &handle_id).await?;
let _driver = driver.lock().await;
let session = state
.broadcast_sessions
.lock()
.await
.remove(&handle_id)
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
session.close();
let adapter = state
.native_broadcast_adapter
.as_deref()
.ok_or_else(|| "native managed broadcast adapter is unavailable".to_string())?;
let result = drive_native_broadcast_actions(&session, adapter).await;
state.broadcast_drivers.lock().await.remove(&handle_id);
result
}
#[tauri::command]
async fn record_sparse_fanout_forward_queue_drop(
state: tauri::State<'_, OpenRtcTauriState>,
capability_key: String,
count: u64,
) -> Result<(), String> {
state
.client()
.record_sparse_fanout_forward_queue_drop(
&required_native_capability_part(capability_key, "capability key")?,
count,
)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn is_peer_connected(
state: tauri::State<'_, OpenRtcTauriState>,
node_id: String,
) -> Result<bool, String> {
state
.client()
.is_connected_str(&node_id)
.await
.map_err(|error| error.to_string())
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct PortableMediaPublicationStart {
publication_id: String,
media_generation: u32,
control: Vec<u8>,
}
fn media_kind_name(kind: openrtc::media::MediaKind) -> &'static str {
match kind {
openrtc::media::MediaKind::Audio => "audio",
openrtc::media::MediaKind::Video => "video",
openrtc::media::MediaKind::Screen => "screen",
openrtc::media::MediaKind::Data => "data",
}
}
fn media_codec_name(codec: openrtc::media::MediaCodec) -> &'static str {
match codec {
openrtc::media::MediaCodec::Opus => "opus",
openrtc::media::MediaCodec::H264 => "h264",
openrtc::media::MediaCodec::Vp8 => "vp8",
openrtc::media::MediaCodec::Vp9 => "vp9",
openrtc::media::MediaCodec::Av1 => "av1",
openrtc::media::MediaCodec::Pcm => "pcm",
openrtc::media::MediaCodec::Opaque => "opaque",
}
}
fn parse_media_publication_id(value: &str) -> Result<openrtc::media::PublicationId, String> {
value
.trim()
.parse()
.map_err(|error: openrtc::media::MediaProtocolError| error.to_string())
}
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn begin_openrtc_media_publication(
state: tauri::State<'_, OpenRtcTauriState>,
publication_id: Option<String>,
kind: String,
codec: String,
clock_rate: u32,
coded_width: Option<u32>,
coded_height: Option<u32>,
channels: Option<u16>,
) -> Result<PortableMediaPublicationStart, String> {
let publication_id = match publication_id.as_deref().map(str::trim) {
Some(value) if !value.is_empty() => parse_media_publication_id(value)?,
_ => openrtc::media::PublicationId::from_bytes(*uuid::Uuid::new_v4().as_bytes()),
};
let kind: openrtc::media::MediaKind = kind
.parse()
.map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
if !matches!(
kind,
openrtc::media::MediaKind::Audio | openrtc::media::MediaKind::Video
) {
return Err("browser media publications must be audio or video".to_string());
}
let codec = codec
.parse()
.map_err(|error: openrtc::media::MediaProtocolError| error.to_string())?;
let publication = openrtc::media::MediaPublicationConfig {
publication_id,
media_generation: 0,
kind,
codec,
clock_rate,
coded_width,
coded_height,
channels,
};
let (publication, control) = state
.portable_media
.lock()
.await
.begin_publication(publication)
.map_err(|error| error.to_string())?;
Ok(PortableMediaPublicationStart {
publication_id: publication.publication_id.to_string(),
media_generation: publication.media_generation,
control,
})
}
#[tauri::command]
#[allow(clippy::too_many_arguments)]
async fn encode_openrtc_media_sample(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<Response, String> {
let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
openrtc::media::ensure_javascript_safe_integer(timestamp_us, "timestamp")
.map_err(|error| error.to_string())?;
let duration_us = parsed_request_header(&request, MEDIA_DURATION_US_HEADER)?;
let keyframe = parsed_request_header(&request, MEDIA_KEYFRAME_HEADER)?;
let discardable = parsed_request_header(&request, MEDIA_DISCARDABLE_HEADER)?;
let payload = request_body_bytes(&request)?;
let encoded = state
.portable_media
.lock()
.await
.encode_sample(
parse_media_publication_id(&publication_id)?,
openrtc::media::EncodedMediaSample {
timestamp_us,
duration_us,
keyframe,
discardable,
payload,
},
)
.map_err(|error| error.to_string())?;
Ok(Response::new(encoded))
}
#[tauri::command]
async fn pause_openrtc_media_publication(
state: tauri::State<'_, OpenRtcTauriState>,
publication_id: String,
) -> Result<(), String> {
state
.portable_media
.lock()
.await
.pause_publication(parse_media_publication_id(&publication_id)?)
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn retire_openrtc_media_publication(
state: tauri::State<'_, OpenRtcTauriState>,
publication_id: String,
) -> Result<(), String> {
state
.portable_media
.lock()
.await
.retire_publication(parse_media_publication_id(&publication_id)?);
Ok(())
}
#[tauri::command]
async fn retire_openrtc_media_receiver(
state: tauri::State<'_, OpenRtcTauriState>,
publication_id: String,
media_generation: u32,
) -> Result<(), String> {
state.portable_media.lock().await.retire_receiver(
parse_media_publication_id(&publication_id)?,
media_generation,
);
Ok(())
}
#[tauri::command]
async fn decode_openrtc_media_chunk(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<Response, String> {
let encoded = request_body_bytes(&request)?;
let chunk = state
.portable_media
.lock()
.await
.decode_chunk(&encoded)
.map_err(|error| error.to_string())?;
openrtc::media::ensure_javascript_safe_integer(chunk.sequence, "sequence")
.and_then(|_| {
openrtc::media::ensure_javascript_safe_integer(chunk.timestamp_us, "timestamp")
})
.map_err(|error| error.to_string())?;
let metadata = serde_json::to_vec(&serde_json::json!({
"publicationId": chunk.publication_id.to_string(),
"mediaGeneration": chunk.media_generation,
"sequence": chunk.sequence,
"timestampUs": chunk.timestamp_us,
"durationUs": chunk.duration_us,
"kind": media_kind_name(chunk.kind),
"codec": media_codec_name(chunk.codec),
"keyframe": chunk.keyframe,
"discardable": chunk.discardable,
}))
.map_err(|error| format!("encode media IPC metadata failed: {error}"))?;
let metadata_len =
u32::try_from(metadata.len()).map_err(|_| "media IPC metadata is too large".to_string())?;
let mut response = Vec::with_capacity(4 + metadata.len() + chunk.payload.len());
response.extend_from_slice(&metadata_len.to_be_bytes());
response.extend_from_slice(&metadata);
response.extend_from_slice(&chunk.payload);
Ok(Response::new(response))
}
#[tauri::command]
async fn decode_openrtc_media_control(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<serde_json::Value, String> {
let encoded = request_body_bytes(&request)?;
let control = state
.portable_media
.lock()
.await
.decode_control(&encoded)
.map_err(|error| error.to_string())?;
Ok(match control {
openrtc::media::MediaControlFrame::Publish {
publication_id,
media_generation,
kind,
codec,
clock_rate,
coded_width,
coded_height,
channels,
} => serde_json::json!({
"type": "publish",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"kind": media_kind_name(kind),
"codec": media_codec_name(codec),
"clockRate": clock_rate,
"codedWidth": coded_width,
"codedHeight": coded_height,
"channels": channels,
}),
openrtc::media::MediaControlFrame::SetEnabled {
publication_id,
media_generation,
enabled,
} => serde_json::json!({
"type": "set-enabled",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"enabled": enabled,
}),
openrtc::media::MediaControlFrame::RequestKeyframe {
publication_id,
media_generation,
} => serde_json::json!({
"type": "request-keyframe",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
}),
openrtc::media::MediaControlFrame::Stop {
publication_id,
media_generation,
reason,
} => serde_json::json!({
"type": "stop",
"publicationId": publication_id.to_string(),
"mediaGeneration": media_generation,
"reason": reason,
}),
})
}
#[tauri::command]
async fn open_peer_bi_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
) -> Result<OpenBiResult, String> {
open_peer_bi_with(app, state, peer_id, timeout_ms, |client, peer_id, timeout_ms| async move {
client.open_peer_bi(&peer_id, timeout_ms).await
})
.await
}
#[tauri::command]
async fn open_peer_bi_transport_only_stream<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
) -> Result<OpenBiResult, String> {
open_peer_bi_with(
app,
state,
peer_id,
timeout_ms,
|client, peer_id, timeout_ms| async move {
client
.open_peer_bi_transport_only(&peer_id, timeout_ms)
.await
.map(|(connection_id, remote_node_id, send, recv)| {
(
connection_id,
remote_node_id,
PeerSendStream::plain(send),
PeerRecvStream::plain(recv),
)
})
},
)
.await
}
#[tauri::command]
async fn open_peer_uni_stream(
state: tauri::State<'_, OpenRtcTauriState>,
peer_id: String,
timeout_ms: Option<u64>,
) -> Result<OpenUniResult, String> {
let peer_id = peer_id.trim().to_string();
if peer_id.is_empty() {
return Err("peerId is required".to_string());
}
let (connection_id, remote_node_id, send) = state
.client()
.open_peer_uni(&peer_id, timeout_ms)
.await
.map_err(|error| format!("open peer uni stream failed: {error}"))?;
let stream_id = uuid::Uuid::new_v4().to_string();
state
.peer_uni_streams
.lock()
.await
.insert(stream_id.clone(), Arc::new(Mutex::new(Some(send))));
Ok(OpenUniResult {
stream_id,
connection_id,
remote_node_id,
})
}
#[tauri::command]
async fn write_peer_bi_stream(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<(), String> {
let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
let bytes = request_body_bytes(&request)?;
let send = {
let streams = state.peer_bi_streams.lock().await;
streams
.get(&stream_id)
.map(|handle| handle.send.clone())
.ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?
};
let result = send
.lock()
.await
.as_mut()
.ok_or_else(|| format!("peer bi stream already closed: {stream_id}"))?
.write_all(&bytes)
.await
.map_err(|error| format!("write peer bi stream failed: {error}"));
result
}
#[tauri::command]
async fn start_peer_bi_stream_read<R: Runtime>(
_app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
channel: Channel<Response>,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let mut streams = state.peer_bi_streams.lock().await;
let handle = streams
.get_mut(&stream_id)
.ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
if handle.read_task.is_some() {
return Ok(());
}
let recv = handle
.recv
.take()
.ok_or_else(|| format!("peer bi stream reader already consumed: {stream_id}"))?;
handle.read_task = Some(spawn_peer_bi_stream_reader(channel, stream_id, recv));
Ok(())
}
#[tauri::command]
async fn cancel_peer_bi_stream_read(
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let mut streams = state.peer_bi_streams.lock().await;
let handle = streams
.get_mut(&stream_id)
.ok_or_else(|| format!("peer bi stream not found: {stream_id}"))?;
handle.recv.take();
if let Some(read_task) = handle.read_task.take() {
read_task.abort();
}
Ok(())
}
#[tauri::command]
async fn close_peer_bi_stream(
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let handle = {
let mut streams = state.peer_bi_streams.lock().await;
streams.remove(&stream_id)
};
if let Some(handle) = handle {
let mut send = handle.send.lock().await;
if let Some(send) = send.take() {
let _ = send.finish();
}
drop(send);
if let Some(read_task) = handle.read_task {
read_task.abort();
}
}
Ok(())
}
#[tauri::command]
async fn write_peer_uni_stream(
state: tauri::State<'_, OpenRtcTauriState>,
request: Request<'_>,
) -> Result<(), String> {
let stream_id = required_request_header(&request, PEER_BI_STREAM_ID_HEADER)?;
let bytes = request_body_bytes(&request)?;
let send = {
let streams = state.peer_uni_streams.lock().await;
streams
.get(&stream_id)
.cloned()
.ok_or_else(|| format!("peer uni stream not found: {stream_id}"))?
};
let result = send
.lock()
.await
.as_mut()
.ok_or_else(|| format!("peer uni stream already closed: {stream_id}"))?
.write_all(&bytes)
.await
.map_err(|error| format!("write peer uni stream failed: {error}"));
result
}
#[tauri::command]
async fn close_peer_uni_stream(
state: tauri::State<'_, OpenRtcTauriState>,
stream_id: String,
) -> Result<(), String> {
let stream_id = stream_id.trim().to_string();
if stream_id.is_empty() {
return Err("streamId is required".to_string());
}
let send = state.peer_uni_streams.lock().await.remove(&stream_id);
if let Some(send) = send {
if let Some(send) = send.lock().await.take() {
let _ = send.finish();
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap as StdHashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex as StdMutex;
use tokio::sync::oneshot;
struct FakeTransportInstaller {
installs: Arc<AtomicUsize>,
}
#[derive(Default)]
struct FakeDeviceKeySigner {
records: StdMutex<StdHashMap<String, String>>,
}
impl DeviceKeySigner for FakeDeviceKeySigner {
fn public_jwk(&self, _app_tag: &str) -> Result<serde_json::Value, String> {
Ok(serde_json::json!({
"kty": "OKP",
"crv": "Ed25519",
"x": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
}))
}
fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>, String> {
assert_eq!(challenge, b"openrtc:v2:test");
Ok(vec![7; 64])
}
fn delete(&self, _app_tag: &str) -> Result<(), String> {
Ok(())
}
fn read_secure_record(&self, app_tag: &str, key: &str) -> Result<Option<String>, String> {
Ok(self
.records
.lock()
.unwrap()
.get(&format!("{app_tag}:{key}"))
.cloned())
}
fn write_secure_record(&self, app_tag: &str, key: &str, value: &str) -> Result<(), String> {
self.records
.lock()
.unwrap()
.insert(format!("{app_tag}:{key}"), value.to_string());
Ok(())
}
fn delete_secure_record(&self, app_tag: &str, key: &str) -> Result<(), String> {
self.records
.lock()
.unwrap()
.remove(&format!("{app_tag}:{key}"));
Ok(())
}
}
impl TransportInstaller for FakeTransportInstaller {
fn id(&self) -> &'static str {
"ble"
}
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
config.ble.as_ref().is_some_and(|ble| ble.enabled)
}
fn install(
&self,
_client: Arc<openrtc::client::Client>,
_config: openrtc::client::TransportConfig,
_context: InstallContext,
) -> InstallFuture {
self.installs.fetch_add(1, Ordering::SeqCst);
Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send + Sync>) })
}
}
#[derive(Debug)]
struct UnavailableTransportInstaller;
impl TransportInstaller for UnavailableTransportInstaller {
fn id(&self) -> &'static str {
"ble"
}
fn is_requested(&self, config: &openrtc::client::TransportConfig) -> bool {
config.ble.as_ref().is_some_and(|ble| ble.enabled)
}
fn install(
&self,
_client: Arc<openrtc::client::Client>,
_config: openrtc::client::TransportConfig,
_context: InstallContext,
) -> InstallFuture {
Box::pin(async { Err("Bluetooth hardware unavailable".to_string()) })
}
}
fn requested_ble_config() -> openrtc::client::TransportConfig {
openrtc::client::TransportConfig {
ble: Some(openrtc::client::BleConfig {
enabled: true,
..openrtc::client::BleConfig::default()
}),
..openrtc::client::TransportConfig::default()
}
}
fn test_config() -> OpenRtcTauriConfig {
OpenRtcTauriConfig {
api_key: format!("pk_test_{}", "a".repeat(40)),
..OpenRtcTauriConfig::default()
}
}
#[test]
fn native_room_architecture_requires_fixed_intent_until_rust_owns_the_gateway_handle() {
let mut registry = NativeCapabilityRegistry::default();
registry
.register_with_architecture(
"room:adaptive".to_string(),
"room".to_string(),
"adaptive".to_string(),
true,
Some(openrtc::native::RoomArchitectureMode::Auto),
)
.expect("room registration");
assert!(
!registry.registrations["room:adaptive"].uses_sparse_fanout(),
"caller-supplied auto state cannot promote sparse fanout",
);
registry
.register_with_architecture(
"room:fixed-sparse".to_string(),
"room".to_string(),
"fixed-sparse".to_string(),
false,
Some(openrtc::native::RoomArchitectureMode::Sparse),
)
.expect("fixed sparse room registration");
assert!(registry.registrations["room:fixed-sparse"].uses_sparse_fanout());
}
#[test]
fn native_room_architecture_rejects_non_room_and_respects_fixed_mesh_intent() {
let mut registry = NativeCapabilityRegistry::default();
assert!(registry
.register_with_architecture(
"space:not-room".to_string(),
"space".to_string(),
"not-room".to_string(),
false,
Some(openrtc::native::RoomArchitectureMode::Auto),
)
.is_err());
registry
.register_with_architecture(
"room:fixed".to_string(),
"room".to_string(),
"fixed".to_string(),
false,
Some(openrtc::native::RoomArchitectureMode::Mesh),
)
.expect("fixed room registration");
assert!(!registry.registrations["room:fixed"].uses_sparse_fanout());
}
#[test]
fn native_capability_registry_isolates_duplicate_peer_channels() {
let mut registry = NativeCapabilityRegistry::default();
registry
.register(
"devices:user-1".to_string(),
"devices".to_string(),
"user-1".to_string(),
false,
)
.expect("devices registration");
registry
.register(
"space:room-1".to_string(),
"space".to_string(),
"room-1".to_string(),
false,
)
.expect("space registration");
registry
.registrations
.get_mut("devices:user-1")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"nodeId": "device-route"
})];
registry
.registrations
.get_mut("space:room-1")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"nodeId": "space-route"
})];
assert_eq!(
registry.capability_keys_for_identity(
Some("shared-peer"),
Some("shared-peer"),
None,
None,
),
vec!["devices:user-1".to_string(), "space:room-1".to_string()]
);
let explicit_channel = openrtc::stream_metadata::ChannelMetadata {
channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
metadata: Some(serde_json::Map::from_iter([(
"openrtcCapability".to_string(),
serde_json::json!("space:room-1"),
)])),
};
assert_eq!(
registry.projected_stream_capability(&explicit_channel),
Some("space:room-1".to_string())
);
let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
metadata: None,
};
assert_eq!(
registry.projected_stream_capability(&unscoped_channel),
None
);
}
#[test]
fn native_capability_registry_aggregates_without_duplicating_root_peers() {
let mut registry = NativeCapabilityRegistry::default();
for (key, kind, id) in [
("devices:user-1", "devices", "user-1"),
("space:room-1", "space", "room-1"),
] {
registry
.register(key.to_string(), kind.to_string(), id.to_string(), false)
.expect("capability registration");
registry.registrations.get_mut(key).unwrap().desired_peers =
vec![serde_json::json!({"deviceId": "shared-peer", "ticket": "ticket-a"})];
}
let (first_revision, peers_json) = registry.aggregate_desired_peers().unwrap();
let peers: Vec<serde_json::Value> = serde_json::from_str(&peers_json).unwrap();
assert_eq!(first_revision, 1);
assert_eq!(peers.len(), 1, "one physical peer must be dialed once");
registry
.register(
"devices:user-1".to_string(),
"devices".to_string(),
"user-1".to_string(),
false,
)
.expect("exact registration is idempotent");
assert!(registry
.register(
"devices:user-1".to_string(),
"space".to_string(),
"room-1".to_string(),
false,
)
.is_err());
}
#[test]
fn native_capability_registry_merges_sparse_roles_deterministically() {
fn aggregate(reverse_registration_order: bool) -> (u64, Vec<serde_json::Value>) {
let mut registry = NativeCapabilityRegistry::default();
let registrations = if reverse_registration_order {
[
("space:backup", "space", "backup"),
("room:active", "room", "active"),
]
} else {
[
("room:active", "room", "active"),
("space:backup", "space", "backup"),
]
};
for (key, kind, id) in registrations {
registry
.register(key.to_string(), kind.to_string(), id.to_string(), false)
.unwrap();
}
registry
.registrations
.get_mut("room:active")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"ticket": "active-ticket",
"topologyRole": "active",
"topologyRevision": 3,
})];
registry
.registrations
.get_mut("space:backup")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"ticket": "backup-ticket",
"topologyRole": "backup",
"topologyRevision": 100,
})];
let (revision, peers_json) = registry.aggregate_desired_peers().unwrap();
(revision, serde_json::from_str(&peers_json).unwrap())
}
let forward = aggregate(false);
let reverse = aggregate(true);
assert_eq!(
forward, reverse,
"registration order is not lifecycle authority"
);
assert_eq!(forward.0, 1);
assert_eq!(forward.1.len(), 1);
assert_eq!(forward.1[0]["ticket"], "active-ticket");
assert_eq!(forward.1[0]["topologyRole"], "active");
assert_eq!(forward.1[0]["topologyRevision"], 1);
}
#[test]
fn native_capability_registry_keeps_physical_peer_until_all_references_withdraw() {
let mut registry = NativeCapabilityRegistry::default();
for (key, kind) in [("room:a", "room"), ("space:b", "space")] {
registry
.register(key.to_string(), kind.to_string(), key.to_string(), false)
.unwrap();
}
registry
.registrations
.get_mut("room:a")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"topologyRole": "backup",
"topologyRevision": 41,
})];
registry
.registrations
.get_mut("space:b")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId": "shared-peer",
"nodeId": "shared-node",
"ticket": "space-ticket",
"topologyRole": "active",
"topologyRevision": 7,
})];
let (_, both_json) = registry.aggregate_desired_peers().unwrap();
let both: Vec<serde_json::Value> = serde_json::from_str(&both_json).unwrap();
assert_eq!(both.len(), 1);
assert_eq!(both[0]["topologyRole"], "active");
assert_eq!(both[0]["ticket"], "space-ticket");
assert_eq!(
registry.registrations["room:a"].desired_peers[0]["topologyRevision"], 41,
"root projection must not rewrite avenue-local lease state",
);
registry
.registrations
.get_mut("space:b")
.unwrap()
.desired_peers
.clear();
let (_, room_only_json) = registry.aggregate_desired_peers().unwrap();
let room_only: Vec<serde_json::Value> = serde_json::from_str(&room_only_json).unwrap();
assert_eq!(
room_only.len(),
1,
"one remaining capability keeps the leg desired"
);
assert_eq!(room_only[0]["topologyRole"], "backup");
registry
.registrations
.get_mut("room:a")
.unwrap()
.desired_peers
.clear();
let (_, empty_json) = registry.aggregate_desired_peers().unwrap();
let empty: Vec<serde_json::Value> = serde_json::from_str(&empty_json).unwrap();
assert!(
empty.is_empty(),
"physical desire ends only after every reference withdraws"
);
}
#[test]
fn native_capability_registry_scopes_peer_data_and_rejects_unscoped_streams() {
let mut registry = NativeCapabilityRegistry::default();
registry
.register(
"devices:user-1".to_string(),
"devices".to_string(),
"user-1".to_string(),
false,
)
.unwrap();
registry
.registrations
.get_mut("devices:user-1")
.unwrap()
.desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
let explicit = openrtc::client::NativePeerDataEvent {
connection_id: "peer-1".to_string(),
remote_node_id: None,
transport: "webrtc".to_string(),
transport_stable_id: 1,
transport_generation: 1,
route_generation: 1,
payload: serde_json::to_vec(&serde_json::json!({
"capability": "devices:user-1",
"kind": "raw",
"payload": [1, 2, 3]
}))
.unwrap(),
};
assert_eq!(
registry.capability_keys_for_peer_data(&explicit),
vec!["devices:user-1".to_string()]
);
let unknown = openrtc::client::NativePeerDataEvent {
payload: serde_json::to_vec(&serde_json::json!({
"capability": "space:unknown",
"kind": "raw"
}))
.unwrap(),
..explicit.clone()
};
assert!(registry.capability_keys_for_peer_data(&unknown).is_empty());
registry
.register(
"space:shared".to_string(),
"space".to_string(),
"shared".to_string(),
false,
)
.unwrap();
registry
.registrations
.get_mut("space:shared")
.unwrap()
.desired_peers = vec![serde_json::json!({"deviceId": "peer-1"})];
let unscoped = openrtc::client::NativePeerDataEvent {
payload: serde_json::to_vec(&serde_json::json!({
"kind": "raw",
"payload": [1, 2, 3]
}))
.unwrap(),
..explicit.clone()
};
assert!(
registry.capability_keys_for_peer_data(&unscoped).is_empty(),
"a shared physical peer cannot fan unscoped data into two avenues",
);
let unscoped_channel = openrtc::stream_metadata::ChannelMetadata {
channel_id: openrtc::stream_metadata::PUBLIC_PEER_STREAM_CHANNEL_ID.to_string(),
metadata: None,
};
assert_eq!(
registry.projected_stream_capability(&unscoped_channel),
None,
"stream ownership must never be inferred from capability count"
);
}
#[test]
fn native_device_proof_uses_host_signer_without_exporting_private_key() {
let state = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
assert_eq!(
state
.native_device_key_signer
.as_deref()
.unwrap()
.offline_assurance("app_native_test"),
openrtc::offline::OfflineAssurance::Software
);
let public = state
.device_public_key("app_native_test")
.expect("public key");
assert_eq!(public["kty"], "OKP");
assert_eq!(public["crv"], "Ed25519");
assert!(public.get("d").is_none());
let signature = state
.sign_device_proof("app_native_test", "openrtc:v2:test")
.expect("signature");
assert_eq!(
base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(signature)
.expect("base64"),
vec![7; 64]
);
}
#[test]
fn offline_support_reports_signer_and_compiled_lan_truth() {
let unavailable = OpenRtcTauriState::new(test_config()).offline_runtime_support();
assert!(!unavailable.provisioning);
assert_eq!(unavailable.local_mesh, cfg!(feature = "transport-lan"));
assert_eq!(
unavailable.reason,
Some("native host device signer is unavailable")
);
let available = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()))
.offline_runtime_support();
assert!(available.provisioning);
assert_eq!(available.local_mesh, cfg!(feature = "transport-lan"));
assert_eq!(
available.reason,
(!cfg!(feature = "transport-lan"))
.then_some("native host was built without transport-lan"),
);
}
#[tokio::test]
async fn public_runtime_status_uses_host_facts_and_keeps_broadcast_unavailable() {
use openrtc::client::CapabilityMaturity;
let unavailable_state = OpenRtcTauriState::new(test_config());
let unavailable = unavailable_state
.project_public_runtime_status(unavailable_state.client().runtime_status().await);
assert_eq!(
unavailable.product_maturity.offline_edge,
CapabilityMaturity::Unavailable
);
assert_eq!(
unavailable.product_maturity.broadcast,
CapabilityMaturity::Unavailable
);
let serialized = serde_json::to_value(&unavailable).expect("serialized runtime status");
assert_eq!(serialized["productMaturity"]["offlineEdge"], "unavailable");
assert_eq!(serialized["productMaturity"]["broadcast"], "unavailable");
let installed_state = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
let installed = installed_state
.project_public_runtime_status(installed_state.client().runtime_status().await);
assert_eq!(
installed.product_maturity.offline_edge,
if cfg!(feature = "transport-lan") {
CapabilityMaturity::Preview
} else {
CapabilityMaturity::SupportOnly
}
);
assert_eq!(
installed.product_maturity.broadcast,
CapabilityMaturity::Unavailable
);
}
#[test]
fn native_certificate_records_round_trip_through_host_secure_storage() {
let state = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
let app_tag = "app_native_test";
let key = "openrtc:v2:device-session:abc:device-1";
assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
state
.write_secure_record(app_tag, key, "{\"token\":\"bound\"}")
.unwrap();
assert_eq!(
state.read_secure_record(app_tag, key).unwrap().as_deref(),
Some("{\"token\":\"bound\"}")
);
state.delete_secure_record(app_tag, key).unwrap();
assert_eq!(state.read_secure_record(app_tag, key).unwrap(), None);
}
#[test]
fn native_device_proof_fails_closed_without_secure_host_signer() {
let state = OpenRtcTauriState::new(test_config());
assert!(state
.device_public_key("app_native_test")
.expect_err("missing signer must fail")
.contains("secure-storage signer"));
}
fn native_transport_install_context() -> InstallContext {
InstallContext {
data_dir: std::env::temp_dir().join("openrtc-tauri-native-transport-tests"),
}
}
fn managed_session_test_result(device_id: &str) -> StartSessionResult {
StartSessionResult {
local_node_id: format!("node-{device_id}"),
ticket_scope: Some("user-device".to_string()),
ticket: None,
presence_started: true,
auto_connect_started: true,
local_device: openrtc::native_device::NativeDeviceIdentity {
device_id: device_id.to_string(),
device_name: "Test Device".to_string(),
created_at_ms: 1,
updated_at_ms: 1,
name_source: None,
system_info: None,
},
}
}
async fn simulate_managed_session_start(
state: Arc<OpenRtcTauriState>,
device_id: String,
pause: Option<(oneshot::Sender<()>, oneshot::Receiver<()>)>,
starts: Arc<AtomicUsize>,
stopped_owners: Arc<Mutex<Vec<usize>>>,
) -> SessionDisposition {
let _start_guard = state.managed_session_start_guard.lock().await;
let client = state.client();
let key = format!("{}:{device_id}", client.app_tag());
let (disposition, active_epoch) = {
let active = state.managed_session.lock().await;
let disposition = managed_session_disposition(
active.as_ref().map(|record| record.key.as_str()),
active
.as_ref()
.is_some_and(|record| Arc::ptr_eq(&record.owner_client, &client)),
active
.as_ref()
.is_some_and(|record| record.result.presence_started),
active
.as_ref()
.is_some_and(|record| record.result.auto_connect_started),
&key,
true,
true,
);
(
disposition,
active.as_ref().map(|record| record.owner_epoch),
)
};
if disposition == SessionDisposition::Reuse {
return disposition;
}
if disposition == SessionDisposition::Replace {
let previous = state
.managed_session
.lock()
.await
.take()
.expect("replacement must have a previous owner");
stopped_owners
.lock()
.await
.push(Arc::as_ptr(&previous.owner_client) as usize);
}
starts.fetch_add(1, Ordering::SeqCst);
if let Some((entered, release)) = pause {
entered.send(()).expect("start observer must be waiting");
release.await.expect("start release must be sent");
}
let owner_epoch = match disposition {
SessionDisposition::Refresh => {
active_epoch.expect("refresh must preserve the active owner epoch")
}
SessionDisposition::Start | SessionDisposition::Replace => {
state.allocate_managed_session_owner_epoch()
}
SessionDisposition::Reuse => unreachable!("reuse returned before startup"),
};
let result = managed_session_test_result(&device_id);
*state.managed_session.lock().await = Some(ManagedSessionRecord {
key,
owner_client: client,
owner_epoch,
result,
});
disposition
}
#[tokio::test]
async fn concurrent_same_key_start_reuses_one_owner_epoch() {
let state = Arc::new(OpenRtcTauriState::new(test_config()));
let starts = Arc::new(AtomicUsize::new(0));
let stopped_owners = Arc::new(Mutex::new(Vec::new()));
let (first_entered_tx, first_entered_rx) = oneshot::channel();
let (first_release_tx, first_release_rx) = oneshot::channel();
let first = tokio::spawn(simulate_managed_session_start(
state.clone(),
"same-device".to_string(),
Some((first_entered_tx, first_release_rx)),
starts.clone(),
stopped_owners.clone(),
));
first_entered_rx.await.expect("first start must pause");
let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
let second_state = state.clone();
let second_starts = starts.clone();
let second_stopped_owners = stopped_owners.clone();
let second = tokio::spawn(async move {
second_attempting_tx.send(()).unwrap();
simulate_managed_session_start(
second_state,
"same-device".to_string(),
None,
second_starts,
second_stopped_owners,
)
.await
});
second_attempting_rx.await.unwrap();
tokio::task::yield_now().await;
assert_eq!(starts.load(Ordering::SeqCst), 1);
first_release_tx.send(()).unwrap();
assert_eq!(
first.await.expect("first start task"),
SessionDisposition::Start
);
assert_eq!(
second.await.expect("second start task"),
SessionDisposition::Reuse
);
assert_eq!(starts.load(Ordering::SeqCst), 1);
let active = state
.managed_session
.lock()
.await
.clone()
.expect("active owner");
assert_eq!(active.owner_epoch, 1);
assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
assert!(stopped_owners.lock().await.is_empty());
}
#[tokio::test]
async fn concurrent_different_key_start_replaces_prior_owner_without_late_overwrite() {
let state = Arc::new(OpenRtcTauriState::new(test_config()));
let starts = Arc::new(AtomicUsize::new(0));
let stopped_owners = Arc::new(Mutex::new(Vec::new()));
let (first_entered_tx, first_entered_rx) = oneshot::channel();
let (first_release_tx, first_release_rx) = oneshot::channel();
let first = tokio::spawn(simulate_managed_session_start(
state.clone(),
"first-device".to_string(),
Some((first_entered_tx, first_release_rx)),
starts.clone(),
stopped_owners.clone(),
));
first_entered_rx.await.expect("first start must pause");
let first_owner = state.client();
let (second_attempting_tx, second_attempting_rx) = oneshot::channel();
let second_state = state.clone();
let second_starts = starts.clone();
let second_stopped_owners = stopped_owners.clone();
let second = tokio::spawn(async move {
second_attempting_tx.send(()).unwrap();
simulate_managed_session_start(
second_state,
"second-device".to_string(),
None,
second_starts,
second_stopped_owners,
)
.await
});
second_attempting_rx.await.unwrap();
tokio::task::yield_now().await;
assert!(Arc::ptr_eq(&first_owner, &state.client()));
first_release_tx.send(()).unwrap();
assert_eq!(
first.await.expect("first start task"),
SessionDisposition::Start
);
assert_eq!(
second.await.expect("replacement start task"),
SessionDisposition::Replace
);
assert_eq!(starts.load(Ordering::SeqCst), 2);
let active = state
.managed_session
.lock()
.await
.clone()
.expect("active owner");
assert_eq!(active.owner_epoch, 2);
assert!(active.key.contains("second-device"));
assert!(Arc::ptr_eq(&active.owner_client, &state.client()));
assert_eq!(
stopped_owners.lock().await.as_slice(),
[Arc::as_ptr(&first_owner) as usize]
);
assert!(Arc::ptr_eq(&active.owner_client, &first_owner));
}
#[test]
fn derives_app_tag_from_api_key_by_default() {
let config = OpenRtcTauriConfig {
api_key: format!("pk_test_{}", "b".repeat(24) + "1234567890abcdef"),
..OpenRtcTauriConfig::default()
};
assert_eq!(config.app_tag().unwrap(), "app_1234567890abcdef");
}
#[tokio::test]
async fn requested_native_transport_without_installer_keeps_base_route_available() {
let state = OpenRtcTauriState::new(test_config());
let client = state.client();
ensure_requested_native_transports(
&state,
&client,
Some(&requested_ble_config()),
&native_transport_install_context(),
)
.await
.expect("an unavailable optional transport must not fail base Iroh startup");
assert!(state.installed_native_transports.lock().await.is_empty());
}
#[tokio::test]
async fn native_transport_installer_is_idempotent_for_one_client() {
let installs = Arc::new(AtomicUsize::new(0));
let state = OpenRtcTauriState::new(test_config()).with_native_transport_installer(
Arc::new(FakeTransportInstaller {
installs: installs.clone(),
}),
);
let client = state.client();
let config = requested_ble_config();
ensure_requested_native_transports(
&state,
&client,
Some(&config),
&native_transport_install_context(),
)
.await
.expect("first install");
ensure_requested_native_transports(
&state,
&client,
Some(&config),
&native_transport_install_context(),
)
.await
.expect("idempotent install");
assert_eq!(installs.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn unavailable_native_transport_keeps_base_route_available() {
let state = OpenRtcTauriState::new(test_config())
.with_native_transport_installer(Arc::new(UnavailableTransportInstaller));
let client = state.client();
ensure_requested_native_transports(
&state,
&client,
Some(&requested_ble_config()),
&native_transport_install_context(),
)
.await
.expect("optional transport installation failure must not fail base Iroh startup");
assert!(state.installed_native_transports.lock().await.is_empty());
}
#[test]
fn invalid_or_missing_api_key_fails_closed() {
assert!(OpenRtcTauriConfig::default().validated_api_key().is_err());
let mut config = OpenRtcTauriConfig::default();
config.api_key = "pk_test_short".to_string();
assert!(config.validated_api_key().is_err());
}
#[cfg(feature = "native-broadcast-moq")]
#[test]
fn draft16_broadcast_feature_installs_one_private_native_adapter() {
let state = OpenRtcTauriState::new(test_config());
assert!(state.native_broadcast_adapter.is_some());
assert!(state.native_broadcast_moq_adapter.is_some());
}
#[test]
fn from_env_uses_api_key_and_ignores_legacy_namespace_selectors() {
let previous_api_key = std::env::var("VITE_OPENRTC_API_KEY").ok();
let previous_project = std::env::var("VITE_OPENRTC_PROJECT_ID").ok();
let previous_app_tag = std::env::var("VITE_PLUTO_OPENRTC_APP_TAG").ok();
let api_key = format!("pk_test_{}", "c".repeat(40));
std::env::set_var("VITE_OPENRTC_API_KEY", &api_key);
std::env::set_var("VITE_OPENRTC_PROJECT_ID", "pluto-rtc-prod");
std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", "app_from_vite_env");
let config = OpenRtcTauriConfig::from_env();
assert_eq!(config.api_key, api_key);
assert_eq!(
config.app_tag().unwrap(),
openrtc::app_tag_from_api_key(&api_key)
);
match previous_api_key {
Some(value) => std::env::set_var("VITE_OPENRTC_API_KEY", value),
None => std::env::remove_var("VITE_OPENRTC_API_KEY"),
}
match previous_project {
Some(value) => std::env::set_var("VITE_OPENRTC_PROJECT_ID", value),
None => std::env::remove_var("VITE_OPENRTC_PROJECT_ID"),
}
match previous_app_tag {
Some(value) => std::env::set_var("VITE_PLUTO_OPENRTC_APP_TAG", value),
None => std::env::remove_var("VITE_PLUTO_OPENRTC_APP_TAG"),
}
}
#[test]
fn native_state_has_one_constructor_owned_app_identity() {
let config = test_config();
let expected = openrtc::app_tag_from_api_key(&config.api_key);
let state = OpenRtcTauriState::new(config);
let first = state.client();
let second = state.client();
assert_eq!(first.app_tag(), expected);
assert!(Arc::ptr_eq(&first, &second));
}
#[test]
fn managed_session_uses_persisted_native_device_id() {
let identity = openrtc::native_device::NativeDeviceIdentity {
device_id: "persisted-native-device".to_string(),
device_name: "Mac".to_string(),
created_at_ms: 1,
updated_at_ms: 1,
name_source: None,
system_info: None,
};
assert_eq!(
managed_session_device_id(Some(" desktop-e2e-native "), &identity),
"persisted-native-device"
);
assert_eq!(
managed_session_device_id(None, &identity),
"persisted-native-device"
);
}
#[test]
fn repeated_managed_session_triggers_converge_on_one_owner() {
let key = "app:persisted-native-device";
let cases = [
("react-remount", true, true, true, true),
("hmr", true, true, true, true),
("auth-refresh", true, true, true, true),
("resume", true, true, true, true),
("alias-change", true, true, true, true),
];
for (label, active_presence, active_auto, requested_presence, requested_auto) in cases {
assert_eq!(
managed_session_disposition(
Some(key),
true,
active_presence,
active_auto,
key,
requested_presence,
requested_auto,
),
SessionDisposition::Reuse,
"{label} must reuse the authoritative tuple"
);
}
assert_eq!(
managed_session_disposition(Some(key), true, false, true, key, true, true),
SessionDisposition::Refresh,
"failed presence startup must retry idempotently"
);
assert_eq!(
managed_session_disposition(
Some(key),
true,
true,
true,
"app:other-native-device",
true,
true,
),
SessionDisposition::Replace,
"a physical native device change must replace the previous lifecycle owner"
);
assert_eq!(
managed_session_disposition(Some(key), false, true, true, key, true, true),
SessionDisposition::Replace,
"the same tuple on a different client epoch must replace the previous owner"
);
}
#[test]
fn user_device_revocation_invalidates_the_cached_managed_ticket() {
assert!(revokes_managed_session("user-device"));
assert!(revokes_managed_session(" user-device "));
assert!(!revokes_managed_session("share:example"));
assert!(!revokes_managed_session(""));
}
#[test]
fn managed_session_metadata_uses_authoritative_device_id() {
let raw = serde_json::json!({
"deviceId": "persisted-native-device",
"assistantDevice": {
"deviceId": "desktop-e2e-native"
}
})
.to_string();
let metadata =
metadata_with_authoritative_device_id(Some(raw), "desktop-e2e-native").unwrap();
let parsed: serde_json::Value = serde_json::from_str(&metadata).unwrap();
assert_eq!(parsed["deviceId"], "desktop-e2e-native");
assert_eq!(parsed["assistantDevice"]["deviceId"], "desktop-e2e-native");
}
}