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 futures::StreamExt;
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::{Emitter, Manager, Runtime, Webview};
use tokio::sync::Mutex;
mod native_assertion;
#[cfg(feature = "native-broadcast-moq")]
mod native_broadcast_moq;
mod native_identity;
pub use native_assertion::{NativeAssertionBridge, NativeAssertionRequest};
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 BROADCAST_HANDLE_ID_HEADER: &str = "x-openrtc-broadcast-handle-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>>,
native_credential: RwLock<Option<(String, Box<dyn Fn() -> Option<String> + Send + Sync>)>>,
}
impl TokenRelayState {
fn token_provider(self: &Arc<Self>) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
let relay = self.clone();
Box::new(move || {
let native = relay.native_credential.read().ok()?;
if let Some((_, provider)) = native.as_ref() {
return provider();
}
drop(native);
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>,
native_source: Option<NativeRosterSource>,
}
#[derive(Debug, Clone)]
enum NativeRosterSource {
Active(String),
Retired,
}
impl NativeCapabilityRegistration {
fn uses_sparse_fanout(&self) -> bool {
self.requested_architecture
.map(|requested| requested == openrtc::native::RoomArchitectureMode::Sparse)
.unwrap_or(self.sparse_fanout)
}
fn owns_native_devices(&self, principal_id: &str) -> bool {
self.avenue_kind == "devices" && self.avenue_id == principal_id
}
}
#[derive(Debug, Default)]
struct NativeCapabilityRegistry {
registrations: HashMap<String, NativeCapabilityRegistration>,
root_desired_revision: u64,
}
impl NativeCapabilityRegistry {
fn native_snapshot(
&mut self,
key: &str,
source_id: &str,
peers: Vec<serde_json::Value>,
) -> Result<Option<(u64, String)>, String> {
let Some(registration) = self.registrations.get_mut(key)
.filter(|registration| matches!(®istration.native_source, Some(NativeRosterSource::Active(id)) if id == source_id)) else {
return Ok(None);
};
registration.desired_peers = peers;
self.aggregate_desired_peers().map(Some)
}
#[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(),
native_source: None,
},
);
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> {
let mut keys = 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(),
)
.into_iter()
.collect::<BTreeSet<_>>();
keys.extend(
snapshot
.scopes
.iter()
.map(|scope| scope.trim())
.filter(|scope| !scope.is_empty() && self.registrations.contains_key(*scope))
.map(str::to_string),
);
keys.into_iter().collect()
}
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,
admitted_scope: Option<&openrtc::session_token::GrantScope>,
) -> 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());
if channel.channel_id.starts_with("share/") {
let session_id = admitted_scope?.0.strip_prefix("share:")?;
let key = format!("ticket:{session_id}");
return (explicit.is_none_or(|claim| claim == key)
&& self.registrations.contains_key(&key))
.then_some(key);
}
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)
}
fn normalize_initial_auto_connect_exclusions(
excluded_peers: Vec<String>,
) -> Result<Vec<String>, String> {
if excluded_peers.len() > 250 {
return Err("native excluded peer list exceeds 250 entries".into());
}
excluded_peers
.into_iter()
.map(|device_id| required_native_capability_part(device_id, "excluded peer device id"))
.collect::<Result<BTreeSet<_>, _>>()
.map(BTreeSet::into_iter)
.map(Iterator::collect)
}
pub struct OpenRtcTauriState {
client: RwLock<Arc<openrtc::client::Client>>,
config: OpenRtcTauriConfig,
token_relay: Arc<TokenRelayState>,
native_identity: Mutex<Option<NativeIdentityContext>>,
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>>>>>,
ticket_meshes: Mutex<HashMap<String, Arc<openrtc::ticket_mesh::TicketMesh>>>,
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>>,
testing_endpoints: Option<(String, String)>,
#[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,
native_identity: Mutex::new(None),
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()),
ticket_meshes: 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,
testing_endpoints: 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
}
pub fn with_testing_endpoints(
mut self,
control_plane: impl Into<String>,
gateway: impl Into<String>,
) -> Self {
self.testing_endpoints = Some((control_plane.into(), gateway.into()));
self
}
pub fn native_control_plane(
&self,
assertions: Arc<dyn openrtc::native::AssertionProvider>,
) -> Result<openrtc::native::ControlPlane, String> {
let signer = self.native_device_key_signer.clone().ok_or_else(|| {
"native control plane requires a host secure-storage signer".to_string()
})?;
let storage = Arc::new(native_identity::HostIdentityStorage(signer));
let control_plane = openrtc::native::ControlPlane::new(
self.config.validated_api_key()?,
assertions,
storage.clone(),
storage,
)
.map_err(|error| error.to_string())?;
let control_plane = if let Some((control_plane_endpoint, gateway_endpoint)) =
self.testing_endpoints.as_ref()
{
control_plane
.with_testing_endpoints(control_plane_endpoint, gateway_endpoint)
.map_err(|error| error.to_string())?
} else {
control_plane
};
Ok(control_plane)
}
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 =
if self.native_device_key_signer.is_some() && self.native_broadcast_adapter.is_some() {
CapabilityMaturity::Preview
} else {
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]>,
) -> Result<native_broadcast_moq::ResolvedNativeBroadcastAccess, 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,
)
.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 {
webview_label: String,
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_configure_native_identity,
openrtc_complete_native_assertion,
openrtc_clear_native_identity,
openrtc_start_native_devices,
openrtc_publish_native_devices,
openrtc_list_native_devices,
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,
issue_ticket_invite,
connect_ticket_invite,
issue_ticket_mesh,
join_ticket_mesh,
ticket_mesh_peers,
ticket_mesh_issuer_node,
close_ticket_mesh,
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,
connect_to_known_device_with_token,
observe_known_device_endpoint,
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,
wait_for_rtc_settled_scope,
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,
pause_openrtc_broadcast_publication,
retire_openrtc_broadcast_publication,
retire_openrtc_broadcast_receiver,
receive_openrtc_broadcast_media,
get_openrtc_broadcast_stats,
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,
finish_peer_bi_stream_send,
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_native_service_error(
registry: &Arc<Mutex<NativeCapabilityRegistry>>,
projection: &Arc<Mutex<Option<NativeProjectionSubscription>>>,
identity_id: &str,
capability_key: &str,
observation: &openrtc::native::NativeServiceErrorObservation,
) -> bool {
let registry = registry.lock().await;
let current = registry.registrations.get(capability_key).is_some_and(|entry| {
matches!(&entry.native_source, Some(NativeRosterSource::Active(id)) if id == identity_id)
&& observation.avenue.kind == "user"
&& entry.owns_native_devices(&observation.avenue.id)
});
if !current { return false; }
send_native_projection(projection, "service-error", vec![capability_key.to_owned()],
&serde_json::json!({
"identityId": identity_id,
"avenue": observation.avenue,
"runtimeInstanceId": observation.runtime_instance_id,
"serviceError": observation.service_error,
})).await;
true
}
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
|| channel.channel_id.starts_with("share/")
}
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 admitted_scope = connection_id.as_deref().and_then(|connection_id| {
match state.client().session_admission(connection_id) {
openrtc::session_token::SessionAdmission::Accepted { scope, .. } => scope,
_ => None,
}
});
let Some(capability_key) = state
.capability_registry
.lock()
.await
.projected_stream_capability(&channel, admitted_scope.as_ref())
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(())
}
struct NativeIdentityContext {
id: String,
webview_label: String,
assertions: Arc<NativeAssertionBridge>,
control_plane: openrtc::native::ControlPlane,
devices: Option<Arc<openrtc::native::Devices>>,
roster: Option<NativeDeviceRoster>,
}
struct NativeDeviceRoster {
capability_key: String,
task: Option<tokio::task::JoinHandle<()>>,
}
fn native_roster_is_active_for_capability(
roster: Option<&NativeDeviceRoster>,
capability_key: &str,
) -> Result<bool, String> {
let Some(roster) = roster else {
return Ok(false);
};
if roster.capability_key != capability_key {
return Err("native devices already have a capability owner".into());
}
Ok(roster.task.as_ref().is_some_and(|task| !task.is_finished()))
}
impl Drop for NativeDeviceRoster {
fn drop(&mut self) {
if let Some(task) = self.task.take() {
task.abort();
}
}
}
impl OpenRtcTauriState {
async fn configure_native_identity(
&self,
webview_label: &str,
session_key: String,
sink: impl Fn(NativeAssertionRequest) -> Result<(), String> + Send + Sync + 'static,
) -> Result<String, String> {
use openrtc::native::AssertionProvider;
let mut current = self.native_identity.lock().await;
if let Some(owner) = current.as_mut() {
if owner.webview_label != webview_label {
return Err("native identity is owned by another webview".into());
}
if owner
.assertions
.session_key()
.map_err(|error| error.to_string())?
.as_deref()
== Some(&session_key)
{
owner
.assertions
.rebind(sink)
.map_err(|error| error.to_string())?;
return Ok(owner.id.clone());
}
}
let assertions = Arc::new(
NativeAssertionBridge::new(session_key, sink).map_err(|error| error.to_string())?,
);
let control_plane = self.native_control_plane(assertions.clone())?;
if let Some(owner) = current.take() {
self.retire_native_identity(owner).await?;
}
let id = uuid::Uuid::new_v4().to_string();
*current = Some(NativeIdentityContext {
id: id.clone(),
webview_label: webview_label.into(),
assertions,
control_plane,
devices: None,
roster: None,
});
Ok(id)
}
}
#[tauri::command]
async fn openrtc_configure_native_identity<R: Runtime>(
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
session_key: String,
requests: Channel<NativeAssertionRequest>,
) -> Result<String, String> {
state
.configure_native_identity(webview.label(), session_key, move |request| {
requests.send(request).map_err(|error| error.to_string())
})
.await
}
#[tauri::command]
async fn openrtc_complete_native_assertion<R: Runtime>(
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
identity_id: String,
request_id: String,
token: Option<String>,
provider_id: Option<String>,
error: Option<String>,
) -> Result<bool, String> {
let current = state.native_identity.lock().await;
let Some(owner) = current.as_ref().filter(|owner| owner.id == identity_id) else {
return Ok(false);
};
if owner.webview_label != webview.label() {
return Err("native identity is owned by another webview".into());
}
let result = match (token, error) {
(Some(token), None) => Ok(openrtc::native::IdentityAssertion { token, provider_id }),
(None, Some(error)) => Err(error),
_ => return Err("provide either an assertion or an error".into()),
};
owner
.assertions
.complete(&request_id, result)
.map_err(|error| error.to_string())
}
impl OpenRtcTauriState {
async fn clear_native_identity(
&self,
webview_label: &str,
identity_id: &str,
) -> Result<bool, String> {
let owner = {
let mut current = self.native_identity.lock().await;
let Some(owner) = current.as_ref().filter(|owner| owner.id == identity_id) else {
return Ok(false);
};
if owner.webview_label != webview_label {
return Err("native identity is owned by another webview".into());
}
current.take().expect("matching identity")
};
self.retire_native_identity(owner).await?;
Ok(true)
}
async fn retire_native_identity(&self, mut owner: NativeIdentityContext) -> Result<(), String> {
drop(owner.roster.take());
owner
.assertions
.retire()
.map_err(|error| error.to_string())?;
{
let mut native = self
.token_relay
.native_credential
.write()
.map_err(|_| "native credential lock poisoned")?;
if native.as_ref().is_some_and(|(id, _)| id == &owner.id) {
*native = None;
self.token_relay.set(None);
}
}
let update = {
let mut registry = self.capability_registry.lock().await;
let mut changed = false;
for registration in registry.registrations.values_mut() {
if matches!(®istration.native_source, Some(NativeRosterSource::Active(id)) if id == &owner.id)
{
registration.native_source = Some(NativeRosterSource::Retired);
registration.desired_peers.clear();
changed = true;
}
}
if changed {
Some(registry.aggregate_desired_peers()?)
} else {
None
}
};
if let Some(devices) = owner.devices {
devices.close().await;
}
if let Some((revision, peers)) = update {
self.client()
.submit_external_desired_peers(revision, &peers)
.await
.map_err(|error| error.to_string())?;
}
Ok(())
}
}
#[tauri::command]
async fn openrtc_clear_native_identity<R: Runtime>(
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
identity_id: String,
) -> Result<bool, String> {
state
.clear_native_identity(webview.label(), &identity_id)
.await
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeDevicesIdentity {
principal_id: String,
device_id: String,
}
#[derive(Debug, Serialize)]
#[serde(untagged)]
enum NativeCommandError {
Service {
#[serde(flatten)]
details: openrtc::service_errors::ServiceError,
message: &'static str,
},
Legacy(String),
}
impl NativeCommandError {
fn from_runtime(error: anyhow::Error) -> Self {
match openrtc::service_errors::ServiceError::from_error(&error) {
Some(details) => Self::Service {
details: details.clone(),
message: "OpenRTC service request denied",
},
None => Self::Legacy(error.to_string()),
}
}
}
impl From<String> for NativeCommandError {
fn from(value: String) -> Self { Self::Legacy(value) }
}
impl From<&str> for NativeCommandError {
fn from(value: &str) -> Self { Self::Legacy(value.to_owned()) }
}
async fn submit_native_device_snapshot(
client: &Arc<openrtc::client::Client>,
registry: &Arc<Mutex<NativeCapabilityRegistry>>,
source_id: &str,
capability_key: &str,
local_device_id: &str,
devices: &[openrtc::signaling::Device],
auto_connect: bool,
) -> Result<bool, String> {
let peers = if auto_connect {
devices
.iter()
.filter(|device| device.device_id != local_device_id)
.map(serde_json::to_value)
.collect::<Result<Vec<_>, _>>()
.map_err(|error| error.to_string())?
} else {
Vec::new()
};
let update = registry
.lock()
.await
.native_snapshot(capability_key, source_id, peers)?;
let Some((revision, aggregate)) = update else {
return Ok(false);
};
client
.submit_external_desired_peers(revision, &aggregate)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn openrtc_list_native_devices<R: Runtime>(
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
identity_id: String,
) -> Result<Vec<openrtc::signaling::Device>, String> {
let devices = {
let current = state.native_identity.lock().await;
let owner = current
.as_ref()
.filter(|owner| owner.id == identity_id)
.ok_or("native identity is no longer current")?;
if owner.webview_label != webview.label() {
return Err("native identity is owned by another webview".into());
}
owner
.devices
.clone()
.ok_or("native devices have not started")?
};
devices
.signaling()
.list_devices(&devices.principal_id, None)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn openrtc_publish_native_devices<R: Runtime>(
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
identity_id: String,
capability_key: String,
device_name: String,
metadata: Option<String>,
excluded_peers: Vec<String>,
auto_connect: bool,
) -> Result<(), NativeCommandError> {
let _start = state.managed_session_start_guard.lock().await;
let key = required_native_capability_part(capability_key, "capability key")?;
let (devices, replay_existing) = {
let current = state.native_identity.lock().await;
let owner = current
.as_ref()
.filter(|owner| owner.id == identity_id)
.ok_or("native identity is no longer current")?;
if owner.webview_label != webview.label() {
return Err("native identity is owned by another webview".into());
}
let replay_existing = native_roster_is_active_for_capability(owner.roster.as_ref(), &key)?;
let devices = owner
.devices
.clone()
.ok_or("native devices have not started")?;
*state
.token_relay
.native_credential
.write()
.map_err(|_| "native credential lock poisoned")? =
Some((identity_id.clone(), devices.identity_credential_provider()));
(devices, replay_existing)
};
if replay_existing {
let snapshot = devices
.signaling()
.list_devices(&devices.principal_id, None)
.await
.map_err(|error| error.to_string())?;
send_native_projection(
&state.native_projection_subscription,
"devices",
vec![key],
&serde_json::json!({ "identityId": identity_id, "devices": snapshot }),
)
.await;
return Ok(());
}
let excluded_peers = normalize_initial_auto_connect_exclusions(excluded_peers)?;
let client = state.client();
for device_id in &excluded_peers {
client.set_auto_connect_excluded(device_id, true);
}
let node_id = ensure_iroh_node(client.clone()).await?;
let local_device_id = client
.get_native_device_identity()
.await
.map_err(|error| error.to_string())?
.device_id;
{
let mut current = state.native_identity.lock().await;
let owner = current
.as_mut()
.filter(|owner| owner.id == identity_id)
.ok_or("native identity changed before presence startup")?;
let mut registry = state.capability_registry.lock().await;
let registration = registry
.registrations
.get_mut(&key)
.ok_or("native capability is not registered")?;
if !registration.owns_native_devices(&devices.principal_id) {
return Err("native devices require their own devices capability registration".into());
}
registration.native_source = Some(NativeRosterSource::Active(identity_id.clone()));
owner.roster = Some(NativeDeviceRoster {
capability_key: key.clone(),
task: None,
});
}
client
.start_external_auto_connect(
format!("{}:native-root", client.app_tag()),
local_device_id.clone(),
)
.await
.map_err(|error| error.to_string())?;
let backend = devices.signaling();
let mut events = backend
.subscribe_devices(&devices.principal_id)
.await
.map_err(|error| error.to_string())?;
let mut service_errors = devices.handle.subscribe_service_errors()
.map_err(|error| error.to_string())?;
backend
.set_excluded_peers(&devices.principal_id, &node_id, &excluded_peers)
.await
.map_err(|error| error.to_string())?;
let ticket = client
.endpoint_ticket_with_token("user-device", 0)
.await
.map_err(|error| error.to_string())?;
backend
.update_presence(
&devices.principal_id,
&node_id,
&ticket,
true,
&device_name,
300_000,
metadata.as_deref(),
)
.await
.map_err(NativeCommandError::from_runtime)?;
let snapshot = backend
.list_devices(&devices.principal_id, None)
.await
.map_err(|error| error.to_string())?;
if !submit_native_device_snapshot(
&client,
&state.capability_registry,
&identity_id,
&key,
&local_device_id,
&snapshot,
auto_connect,
)
.await?
{
return Err("native device roster was superseded".into());
}
let mut current = state.native_identity.lock().await;
let Some(owner) = current.as_mut().filter(|owner| owner.id == identity_id) else {
drop(current);
devices.close().await;
return Err("native identity changed during presence startup".into());
};
send_native_projection(
&state.native_projection_subscription,
"devices",
vec![key.clone()],
&serde_json::json!({ "identityId": identity_id, "devices": snapshot }),
)
.await;
let registry = state.capability_registry.clone();
let projection = state.native_projection_subscription.clone();
let task_key = key.clone();
let task = tokio::spawn(async move {
loop {
let event = tokio::select! {
biased;
observation = service_errors.recv() => {
match observation {
Ok(observation) => {
if devices.is_closed() || !forward_native_service_error(
®istry, &projection, &identity_id, &task_key, &observation,
).await { break; }
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
continue;
}
event = events.next() => match event { Some(event) => event, None => break },
};
if let Err(error) = event {
eprintln!("[openrtc-tauri] native roster notification failed: {error}");
}
let snapshot = match backend.list_devices(&devices.principal_id, None).await {
Ok(snapshot) => snapshot,
Err(error) => {
eprintln!("[openrtc-tauri] native roster read failed: {error}");
continue;
}
};
match submit_native_device_snapshot(
&client,
®istry,
&identity_id,
&task_key,
&local_device_id,
&snapshot,
auto_connect,
)
.await
{
Ok(true) => {
send_native_projection(
&projection,
"devices",
vec![task_key.clone()],
&serde_json::json!({ "identityId": identity_id, "devices": snapshot }),
)
.await
}
Ok(false) => break,
Err(error) => {
eprintln!("[openrtc-tauri] native desired-peer input failed: {error}");
break;
}
}
}
});
owner.roster = Some(NativeDeviceRoster {
capability_key: key,
task: Some(task),
});
Ok(())
}
#[tauri::command]
async fn openrtc_start_native_devices<R: Runtime>(
app: tauri::AppHandle<R>,
webview: Webview<R>,
state: tauri::State<'_, OpenRtcTauriState>,
identity_id: String,
device_name: Option<String>,
) -> Result<NativeDevicesIdentity, NativeCommandError> {
let _start = state.managed_session_start_guard.lock().await;
let device = state
.client()
.init_native_device_identity(app_data_dir(&app, &state)?, device_name.as_deref())
.await
.map_err(|error| error.to_string())?;
let control_plane = {
let current = state.native_identity.lock().await;
let owner = current
.as_ref()
.filter(|owner| owner.id == identity_id)
.ok_or("native identity is no longer current")?;
if owner.webview_label != webview.label() {
return Err("native identity is owned by another webview".into());
}
if let Some(devices) = &owner.devices {
if !devices.is_closed() {
return Ok(NativeDevicesIdentity {
principal_id: devices.principal_id.clone(),
device_id: device.device_id,
});
}
}
owner.control_plane.clone()
};
let devices = Arc::new(
control_plane
.devices(
device.device_id.clone(),
std::env::consts::OS,
openrtc::native::DeviceOptions {
device_profile: Some(openrtc::native::DeviceProfile {
name: device_name,
platform: Some(std::env::consts::OS.into()),
}),
..Default::default()
},
)
.await
.map_err(NativeCommandError::from_runtime)?,
);
let mut current = state.native_identity.lock().await;
let Some(owner) = current.as_mut().filter(|owner| owner.id == identity_id) else {
drop(current);
devices.close().await;
return Err("native identity changed during device enrollment".into());
};
let result = NativeDevicesIdentity {
principal_id: devices.principal_id.clone(),
device_id: device.device_id,
};
owner.devices = Some(devices);
Ok(result)
}
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 issue_ticket_invite(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
max_peers: u32,
mesh: bool,
) -> Result<String, String> {
ensure_iroh_node(state.client()).await?;
let invite = if mesh {
state.client().issue_ticket_invite(
&id, openrtc::ticket_mesh::TicketMeshOptions { max_peers },
).await
} else {
state.client().issue_direct_ticket_invite(&id, max_peers).await
};
invite
.map(|invite| invite.encode())
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn connect_ticket_invite(
state: tauri::State<'_, OpenRtcTauriState>,
invitation: String,
) -> Result<openrtc::client::ManagedConnectResult, String> {
ensure_iroh_node(state.client()).await?;
state
.client()
.connect_ticket_invite(&invitation)
.await
.map_err(|error| error.to_string())
}
#[tauri::command]
async fn issue_ticket_mesh<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
max_peers: u32,
) -> Result<String, String> {
let client = state.client();
ensure_iroh_node(client.clone()).await?;
let mut meshes = state.ticket_meshes.lock().await;
if meshes.contains_key(&id) {
return Err("ticket mesh already has an owner on this runtime".into());
}
let mesh = openrtc::ticket_mesh::TicketMesh::issue(
client,
&id,
openrtc::ticket_mesh::TicketMeshOptions { max_peers },
)
.await
.map_err(|error| error.to_string())?;
let invite = mesh.invite();
forward_ticket_mesh_invites(app, id.clone(), mesh.watch_invite());
meshes.insert(id, Arc::new(mesh));
Ok(invite)
}
#[tauri::command]
async fn join_ticket_mesh<R: Runtime>(
app: tauri::AppHandle<R>,
state: tauri::State<'_, OpenRtcTauriState>,
invitation: String,
) -> Result<String, String> {
let invite = openrtc::ticket_mesh::TicketInvite::parse(&invitation)
.map_err(|error| error.to_string())?;
let id = invite.id().to_string();
let client = state.client();
ensure_iroh_node(client.clone()).await?;
let mut meshes = state.ticket_meshes.lock().await;
if meshes.contains_key(&id) {
return Err("ticket mesh already has an owner on this runtime".into());
}
let mesh = openrtc::ticket_mesh::TicketMesh::join(client, &invitation)
.await
.map_err(|error| error.to_string())?;
forward_ticket_mesh_invites(app, id.clone(), mesh.watch_invite());
meshes.insert(id.clone(), Arc::new(mesh));
Ok(id)
}
#[derive(Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct TicketMeshInviteEvent {
id: String,
invitation: String,
}
fn forward_ticket_mesh_invites<R: Runtime>(
app: tauri::AppHandle<R>,
id: String,
mut invites: tokio::sync::watch::Receiver<String>,
) {
tauri::async_runtime::spawn(async move {
while invites.changed().await.is_ok() {
let invitation = invites.borrow_and_update().clone();
let _ = app.emit("openrtc://ticket-mesh-invite", TicketMeshInviteEvent {
id: id.clone(),
invitation,
});
}
});
}
#[tauri::command]
async fn ticket_mesh_peers(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<Vec<openrtc::client::PeerSessionSnapshot>, String> {
let mesh = state.ticket_meshes.lock().await.get(&id).cloned()
.ok_or_else(|| "ticket mesh is not active".to_string())?;
Ok(mesh.watch_peers().borrow().clone())
}
#[tauri::command]
async fn ticket_mesh_issuer_node(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<String, String> {
let mesh = state.ticket_meshes.lock().await.get(&id).cloned()
.ok_or_else(|| "ticket mesh is not active".to_string())?;
mesh.issuer_node().map_err(|error| error.to_string())
}
#[tauri::command]
async fn close_ticket_mesh(
state: tauri::State<'_, OpenRtcTauriState>,
id: String,
) -> Result<(), String> {
let mesh = state.ticket_meshes.lock().await.get(&id).cloned();
if let Some(mesh) = mesh {
mesh.close().await;
let mut meshes = state.ticket_meshes.lock().await;
if meshes
.get(&id)
.is_some_and(|current| Arc::ptr_eq(current, &mesh))
{
meshes.remove(&id);
}
}
Ok(())
}
#[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 native_devices = {
let mut current = state.native_identity.lock().await;
match current.as_mut().filter(|owner| {
owner
.roster
.as_ref()
.is_some_and(|roster| roster.capability_key == capability_key)
}) {
Some(owner) => {
drop(owner.roster.take());
owner
.devices
.take()
.map(|devices| (owner.id.clone(), devices))
}
None => None,
}
};
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 let Some((identity_id, devices)) = native_devices {
{
let mut native = state
.token_relay
.native_credential
.write()
.map_err(|_| "native credential lock poisoned")?;
if native.as_ref().is_some_and(|(id, _)| id == &identity_id) {
*native = None;
}
}
devices.close().await;
}
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 registration.native_source.is_some() {
return Err("desired peers for this capability are owned by the native gateway".into());
}
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 connect_to_known_device_with_token(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
token: String,
scope: String,
max_connections: u32,
expires_at_ms: Option<u64>,
timeout_ms: Option<u64>,
) -> Result<openrtc::client::ManagedConnectResult, String> {
let client = state.client();
let connect = client.connect_known_device_with_token(
&device_id,
&token,
&scope,
max_connections,
expires_at_ms,
timeout_ms.map(|value| value.min(10_000)),
);
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_known_device_with_token timeout after {timeout_ms}ms")
})?
.map_err(|error| error.to_string())
}
None => connect.await.map_err(|error| error.to_string()),
}
}
#[tauri::command]
async fn observe_known_device_endpoint(
state: tauri::State<'_, OpenRtcTauriState>,
device_id: String,
endpoint_ticket: String,
) -> Result<(), String> {
state
.client()
.observe_known_device_endpoint(&device_id, &endpoint_ticket)
.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> {
let client = state.client();
client.set_auto_connect_excluded(&device_id, excluded);
client.wake_native_external_auto_connect().await;
let devices = state
.native_identity
.lock()
.await
.as_ref()
.and_then(|owner| owner.devices.clone());
if let (Some(devices), Some(node_id)) = (devices, client.current_node_id().await) {
devices
.signaling()
.set_excluded_peers(
&devices.principal_id,
&node_id,
&client.current_excluded_peers_snapshot(),
)
.await
.map_err(|error| error.to_string())?;
}
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 wait_for_rtc_settled_scope(
state: tauri::State<'_, OpenRtcTauriState>,
scope: String,
timeout_ms: Option<u64>,
) -> Result<Option<openrtc::client::PeerSessionSnapshot>, String> {
Ok(state.client().wait_for_settled_scope(&scope, 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<R: Runtime>(
webview: Webview<R>,
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;
install_native_projection(
&mut subscription,
webview.label().to_string(),
request_id,
channel,
)
}
fn install_native_projection(
subscription: &mut Option<NativeProjectionSubscription>,
webview_label: String,
request_id: String,
channel: Channel<NativeProjectionEvent>,
) -> Result<String, String> {
if let Some(current) = subscription.as_ref() {
if current.webview_label == webview_label {
*subscription = Some(NativeProjectionSubscription {
webview_label,
request_id: request_id.clone(),
channel,
});
return Ok(request_id);
}
return Err("OpenRTC native projection subscription already has a process owner".into());
}
*subscription = Some(NativeProjectionSubscription {
webview_label,
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,
source_slots: Vec<String>,
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: Option<Vec<u8>>,
publisher_signer_handle: Option<String>,
now_ms: u64,
) -> Result<NativeBroadcastSessionDescriptor, 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
};
#[cfg(feature = "native-broadcast-moq")]
let resolved_access = if state.native_broadcast_moq_adapter.is_some() {
Some(
state
.resolve_native_broadcast_access(&grant_token, publication_verifying_key.as_ref())
.await?,
)
} else {
None
};
#[cfg(feature = "native-broadcast-moq")]
let issuer_public_key = match (issuer_public_key, resolved_access.as_ref()) {
(Some(provided), Some(resolved)) if provided.as_slice() != resolved.issuer_public_key => {
return Err("broadcast issuer key mismatch".to_string());
}
(Some(provided), _) => provided,
(None, Some(resolved)) => resolved.issuer_public_key.to_vec(),
(None, None) => return Err("broadcast issuer key is required".to_string()),
};
#[cfg(not(feature = "native-broadcast-moq"))]
let issuer_public_key =
issuer_public_key.ok_or_else(|| "broadcast issuer key is required".to_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())?;
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();
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(),
resolved_access.map(|resolved| resolved.relay),
) {
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(),
source_slots: session.source_slots(),
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>,
request: Request<'_>,
) -> Result<(), String> {
let handle_id = required_request_header(&request, BROADCAST_HANDLE_ID_HEADER)?;
let publication_id = required_request_header(&request, MEDIA_PUBLICATION_ID_HEADER)?;
let timestamp_us = parsed_request_header(&request, MEDIA_TIMESTAMP_US_HEADER)?;
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 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
}
#[tauri::command]
async fn pause_openrtc_broadcast_publication(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
publication_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
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
session
.pause_publication(parse_media_publication_id(&publication_id)?)
.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 retire_openrtc_broadcast_publication(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
publication_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
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
session
.retire_publication(parse_media_publication_id(&publication_id)?)
.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 retire_openrtc_broadcast_receiver(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
publication_id: String,
media_generation: u32,
) -> 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
.retire_receiver(
parse_media_publication_id(&publication_id)?,
media_generation,
)
.map_err(|error| error.to_string())
}
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<openrtc::broadcast::AcceptedBroadcastMedia>, 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 media = if let Some(source_slot) = received.source_slot.as_deref() {
session
.accept_media_for_source(source_slot, &received.object)
.map_err(|error| error.to_string())?
} else {
session
.accept_media(&received.object)
.map_err(|error| error.to_string())?
};
let Some(media) = media else { continue };
delivered_bytes = delivered_bytes.saturating_add(object_len);
if let Some(encoded) = media.media.as_deref() {
let chunk = openrtc::media::EncodedMediaChunk::decode(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())?;
}
accepted.push(media);
}
#[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 get_openrtc_broadcast_stats(
state: tauri::State<'_, OpenRtcTauriState>,
handle_id: String,
) -> Result<openrtc::broadcast::BroadcastStats, String> {
let session = state
.broadcast_sessions
.lock()
.await
.get(&handle_id)
.cloned()
.ok_or_else(|| "broadcast session is unavailable".to_string())?;
Ok(session.stats())
}
#[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 finish_peer_bi_stream_send(
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 = {
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}"))?
};
if let Some(send) = send.lock().await.take() {
send.finish()
.map_err(|error| format!("finish peer bi stream send failed: {error}"))?;
}
Ok(())
}
#[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;
#[test]
fn initial_auto_connect_exclusions_are_validated_and_deduplicated() {
assert_eq!(
normalize_initial_auto_connect_exclusions(vec![
" device-b ".into(),
"device-a".into(),
"device-b".into(),
])
.unwrap(),
vec!["device-a", "device-b"],
);
assert!(normalize_initial_auto_connect_exclusions(vec!["".into()]).is_err());
assert!(normalize_initial_auto_connect_exclusions(vec!["x".into(); 251]).is_err());
}
#[test]
fn native_startup_denials_keep_structured_metadata_and_legacy_strings() {
for code in openrtc::service_errors::SERVICE_ERROR_CODES {
let error = NativeCommandError::Service {
message: "OpenRTC service request denied",
details: openrtc::service_errors::ServiceError::from_value(&serde_json::json!({
"code": code, "retryable": false, "scope": "app", "operation": "gateway.grant.issue",
"requestId": "startup-1", "retryAfterMs": 1200, "resetAt": 1800000000000_u64,
"providerCost": 99,
})).unwrap(),
};
let value = serde_json::to_value(error).unwrap();
assert_eq!(value["code"], *code);
assert_eq!(value["retryable"], false);
assert_eq!(value["scope"], "app");
assert_eq!(value["requestId"], "startup-1");
assert_eq!(value["retryAfterMs"], 1200);
assert_eq!(value["resetAt"], 1800000000000_u64);
assert_eq!(value["message"], "OpenRTC service request denied");
assert!(value.get("providerCost").is_none());
}
let legacy = NativeCommandError::from_runtime(anyhow::anyhow!("host unavailable"));
assert_eq!(serde_json::to_value(legacy).unwrap(), "host unavailable");
}
#[cfg(feature = "testing-endpoints")]
#[tokio::test]
async fn real_http_denial_reaches_native_command_encoder() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
struct DenialSigner;
impl openrtc::native::DeviceSigner for DenialSigner {
fn public_jwk(&self, _: &str) -> anyhow::Result<serde_json::Value> {
Ok(serde_json::json!({ "kty": "OKP", "crv": "Ed25519",
"x": "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA" }))
}
fn sign(&self, _: &str, _: &[u8]) -> anyhow::Result<Vec<u8>> {
Ok(vec![0; 64])
}
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut stream, _) = tokio::time::timeout(std::time::Duration::from_secs(5), listener.accept())
.await.unwrap().unwrap();
let mut bytes = Vec::new();
let mut buffer = [0_u8; 4096];
loop {
let read = tokio::time::timeout(std::time::Duration::from_secs(5), stream.read(&mut buffer))
.await.unwrap().unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
assert!(bytes.len() < 64 * 1024);
if let Some(end) = bytes.windows(4).position(|part| part == b"\r\n\r\n") {
let length = String::from_utf8_lossy(&bytes[..end]).lines().find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length").then(|| value.trim().parse::<usize>().unwrap())
}).unwrap();
assert!(length < 32 * 1024);
if bytes.len() >= end + 4 + length { break; }
}
}
assert!(String::from_utf8_lossy(&bytes).starts_with("POST /v2/capabilities "));
let body = r#"{"error":"private provider response","code":"app-budget-exhausted","retryable":false,"scope":"app","operation":"capability.issue","requestId":"local-denial","retryAfterMs":1200,"providerCost":99}"#;
stream.write_all(format!(
"HTTP/1.1 429 Too Many Requests\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", body.len(), body,
).as_bytes()).await.unwrap();
});
let control = openrtc::native::ControlPlane::anonymous(
"pk_live_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
Arc::new(DenialSigner),
).unwrap().with_testing_endpoints(format!("http://{address}"), "http://127.0.0.1:1").unwrap();
let result = tokio::time::timeout(std::time::Duration::from_secs(5),
control.join_space("local-room", "desktop", openrtc::native::CapabilityOptions::default()),
).await.expect("local denial must settle without retry");
let error = match result { Err(error) => error, Ok(_) => panic!("denied capability admitted") };
server.await.unwrap();
let value = serde_json::to_value(NativeCommandError::from_runtime(error.context("private host context"))).unwrap();
assert_eq!(value["code"], "app-budget-exhausted");
assert_eq!(value["retryable"], false);
assert_eq!(value["operation"], "capability.issue");
assert_eq!(value["requestId"], "local-denial");
assert_eq!(value["retryAfterMs"], 1200);
assert_eq!(value["message"], "OpenRTC service request denied");
assert!(!value.to_string().contains("private"));
assert!(value.get("providerCost").is_none());
}
#[test]
fn native_projection_rejects_competing_owner_without_replacing_channel() {
let mut slot = None;
let first = Channel::new(|_| Ok(()));
let first_id = first.id();
install_native_projection(&mut slot, "main".into(), "first".into(), first).unwrap();
let error = install_native_projection(
&mut slot,
"settings".into(),
"second".into(),
Channel::new(|_| Ok(())),
)
.unwrap_err();
assert!(error.contains("already has a process owner"));
assert_eq!(slot.as_ref().unwrap().request_id, "first");
assert_eq!(slot.as_ref().unwrap().channel.id(), first_id);
let refreshed = Channel::new(|_| Ok(()));
let refreshed_id = refreshed.id();
install_native_projection(&mut slot, "main".into(), "first".into(), refreshed).unwrap();
assert_eq!(slot.as_ref().unwrap().channel.id(), refreshed_id);
slot = None;
install_native_projection(
&mut slot,
"settings".into(),
"second".into(),
Channel::new(|_| Ok(())),
)
.unwrap();
assert_eq!(slot.as_ref().unwrap().request_id, "second");
assert!(install_native_projection(
&mut slot,
"main".into(),
"first".into(),
Channel::new(|_| Ok(())),
)
.is_err());
assert_eq!(slot.as_ref().unwrap().request_id, "second");
}
#[test]
fn native_projection_same_webview_reclaims_after_reload() {
let mut slot = None;
install_native_projection(
&mut slot,
"main".into(),
"before-reload".into(),
Channel::new(|_| Ok(())),
)
.unwrap();
install_native_projection(
&mut slot,
"main".into(),
"after-reload".into(),
Channel::new(|_| Ok(())),
)
.unwrap();
assert_eq!(slot.as_ref().unwrap().webview_label, "main");
assert_eq!(slot.as_ref().unwrap().request_id, "after-reload");
}
#[tokio::test]
async fn native_service_errors_reject_retired_identity_and_wrong_avenue() {
let registry = Arc::new(Mutex::new(NativeCapabilityRegistry::default()));
{
let mut state = registry.lock().await;
state.register("devices:p".into(), "devices".into(), "p".into(), false).unwrap();
state.registrations.get_mut("devices:p").unwrap().native_source =
Some(NativeRosterSource::Active("identity-new".into()));
}
let received = Arc::new(StdMutex::new(Vec::<serde_json::Value>::new()));
let sink = received.clone();
let channel = Channel::new(move |body| {
if let tauri::ipc::InvokeResponseBody::Json(json) = body {
sink.lock().unwrap().push(serde_json::from_str(&json).unwrap());
}
Ok(())
});
let projection = Arc::new(Mutex::new(Some(NativeProjectionSubscription {
webview_label: "main".into(), request_id: "current-request".into(), channel,
})));
let mut observation = openrtc::native::NativeServiceErrorObservation {
avenue: serde_json::from_value(serde_json::json!({ "kind": "user", "id": "p" })).unwrap(),
runtime_instance_id: "runtime:1".into(),
service_error: openrtc::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-budget-exhausted", "retryable": false, "scope": "app", "providerCost": 99,
})).unwrap(),
};
assert!(!forward_native_service_error(®istry, &projection, "identity-old", "devices:p", &observation).await);
observation.avenue.id = "another-principal".into();
assert!(!forward_native_service_error(®istry, &projection, "identity-new", "devices:p", &observation).await);
observation.avenue.id = "p".into();
assert!(forward_native_service_error(®istry, &projection, "identity-new", "devices:p", &observation).await);
let events = received.lock().unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0]["kind"], "service-error");
assert_eq!(events[0]["requestId"], "current-request");
assert_eq!(events[0]["capabilityKeys"], serde_json::json!(["devices:p"]));
assert_eq!(events[0]["payload"]["identityId"], "identity-new");
assert_eq!(events[0]["payload"]["serviceError"]["code"], "app-budget-exhausted");
assert!(events[0]["payload"]["serviceError"].get("providerCost").is_none());
drop(events);
registry.lock().await.registrations.get_mut("devices:p").unwrap().native_source =
Some(NativeRosterSource::Retired);
assert!(!forward_native_service_error(®istry, &projection, "identity-new", "devices:p", &observation).await);
registry.lock().await.registrations.remove("devices:p");
assert!(!forward_native_service_error(®istry, &projection, "identity-new", "devices:p", &observation).await);
assert_eq!(received.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn active_native_roster_is_replayed_for_the_same_capability_after_reload() {
let (_keep_alive, wait) = oneshot::channel::<()>();
let roster = NativeDeviceRoster {
capability_key: "devices:user-1".into(),
task: Some(tokio::spawn(async move {
let _ = wait.await;
})),
};
assert!(native_roster_is_active_for_capability(Some(&roster), "devices:user-1",).unwrap());
assert!(
native_roster_is_active_for_capability(Some(&roster), "devices:user-2",)
.unwrap_err()
.contains("already have a capability owner")
);
}
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 public_devices_capability_owns_the_matching_native_user_avenue() {
let registration = NativeCapabilityRegistration {
avenue_kind: "devices".to_string(),
avenue_id: "principal-1".to_string(),
sparse_fanout: false,
requested_architecture: None,
desired_revision: 0,
desired_peers: Vec::new(),
native_source: None,
};
assert!(registration.owns_native_devices("principal-1"));
assert!(!registration.owns_native_devices("principal-2"));
let mut wrong_kind = registration.clone();
wrong_kind.avenue_kind = "user".to_string();
assert!(!wrong_kind.owns_native_devices("principal-1"));
}
#[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, None),
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),
None
);
}
#[test]
fn native_capability_registry_routes_state_by_authenticated_scope() {
let mut registry = NativeCapabilityRegistry::default();
registry
.register(
"ticket:share-1".to_string(),
"ticket".to_string(),
"share-1".to_string(),
false,
)
.expect("ticket registration");
registry
.register(
"space:other".to_string(),
"space".to_string(),
"other".to_string(),
false,
)
.expect("unrelated registration");
let snapshot = openrtc::client::StateSnapshot {
connection_id: "mesh-connection".to_string(),
device_id: Some("mesh-peer".to_string()),
device_id_hint: None,
remote_node_id: Some("mesh-node".to_string()),
state: "connected".to_string(),
transport_state: "connected".to_string(),
protocol_state: "routable".to_string(),
routable: true,
logical_session_terminal: false,
readiness_state: openrtc::client::ReadinessState::Routable,
readiness_reason: "settled".to_string(),
transport_generation: 1,
route_generation: 1,
active_transport_stable_id: Some(1),
active_transport: "iroh".to_string(),
parallel_transport: None,
replacement_in_progress: false,
last_lifecycle_transition_at_ms: 1,
transition_count: 1,
connecting_transition_count: 1,
replacement_count: 0,
retire_count: 0,
last_disconnect_reason: None,
last_reconnect_reason: None,
scopes: vec![
" unregistered:scope ".to_string(),
" ticket:share-1 ".to_string(),
"ticket:share-1".to_string(),
],
error: None,
created_at: 1,
updated_at: 1,
};
assert_eq!(
registry.capability_keys_for_state(&snapshot),
vec!["ticket:share-1".to_string()],
"the registered authenticated scope routes the state without a desired-peer roster",
);
}
#[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 ticket_share_stream_projection_follows_rust_admission_scope() {
let mut registry = NativeCapabilityRegistry::default();
registry.register(
"ticket:share_host_123".to_string(),
"ticket".to_string(),
"share_host_123".to_string(),
false,
).unwrap();
let admitted = openrtc::session_token::GrantScope("share:share_host_123".to_string());
let wrong_scope = openrtc::session_token::GrantScope("share:other".to_string());
let channel = openrtc::stream_metadata::ChannelMetadata {
channel_id: "share/explicit-control".to_string(),
metadata: Some(serde_json::Map::from_iter([(
"openrtcCapability".to_string(),
serde_json::json!("ticket:share_host_123"),
)])),
};
assert!(is_openrtc_projected_channel(&channel));
assert_eq!(registry.projected_stream_capability(&channel, Some(&admitted)),
Some("ticket:share_host_123".to_string()));
assert_eq!(registry.projected_stream_capability(&channel, Some(&wrong_scope)), None);
assert_eq!(registry.projected_stream_capability(&channel, None), None);
let mut untagged = channel.clone();
untagged.metadata = None;
assert!(is_openrtc_projected_channel(&untagged));
assert_eq!(registry.projected_stream_capability(&untagged, Some(&admitted)),
Some("ticket:share_host_123".to_string()));
let mut wrong_owner = channel;
wrong_owner.metadata = Some(serde_json::Map::from_iter([(
"openrtcCapability".to_string(),
serde_json::json!("devices:user-1"),
)]));
assert_eq!(registry.projected_stream_capability(&wrong_owner, Some(&admitted)), None);
}
#[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),
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]
);
}
#[tokio::test]
async fn native_identity_reload_preserves_owner_and_principal_change_retires_it() {
use openrtc::native::AssertionProvider;
let state = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
let first = state
.configure_native_identity("main", "first-user".into(), |_| Ok(()))
.await
.unwrap();
let provider = state
.native_identity
.lock()
.await
.as_ref()
.unwrap()
.assertions
.clone();
let repeated = state
.configure_native_identity("main", "first-user".into(), |_| Ok(()))
.await
.unwrap();
assert_eq!(first, repeated);
assert!(Arc::ptr_eq(
&provider,
&state
.native_identity
.lock()
.await
.as_ref()
.unwrap()
.assertions
));
assert_eq!(provider.identity_epoch(), 0);
assert!(state
.configure_native_identity("other-window", "second-user".into(), |_| Ok(()))
.await
.is_err());
assert_eq!(
provider.identity_epoch(),
0,
"competing window cannot retire the owner"
);
let second = state
.configure_native_identity("main", "second-user".into(), |_| Ok(()))
.await
.unwrap();
assert_ne!(first, second);
assert_eq!(provider.identity_epoch(), 1);
assert!(provider.session_key().unwrap().is_none());
let owner = state.native_identity.lock().await;
assert_eq!(
owner
.as_ref()
.unwrap()
.assertions
.session_key()
.unwrap()
.as_deref(),
Some("second-user")
);
let next_provider = owner.as_ref().unwrap().assertions.clone();
drop(owner);
assert!(!state.clear_native_identity("main", &first).await.unwrap());
assert!(state
.clear_native_identity("other-window", &second)
.await
.is_err());
assert_eq!(next_provider.identity_epoch(), 0);
assert!(state.clear_native_identity("main", &second).await.unwrap());
assert_eq!(next_provider.identity_epoch(), 1);
assert!(state.native_identity.lock().await.is_none());
}
#[tokio::test]
async fn native_roster_feeds_root_and_fences_replaced_source() {
let state = OpenRtcTauriState::new(test_config())
.with_native_device_key_signer(Arc::new(FakeDeviceKeySigner::default()));
let old_source = state
.configure_native_identity("main", "first-user".into(), |_| Ok(()))
.await
.unwrap();
let client = state.client();
client
.start_external_auto_connect(
format!("{}:native-root", client.app_tag()),
"local".into(),
)
.await
.unwrap();
{
let mut registry = state.capability_registry.lock().await;
registry
.register("devices".into(), "user".into(), "principal".into(), false)
.unwrap();
registry
.register("space".into(), "space".into(), "shared-space".into(), false)
.unwrap();
registry
.registrations
.get_mut("devices")
.unwrap()
.native_source = Some(NativeRosterSource::Active(old_source.clone()));
registry
.registrations
.get_mut("space")
.unwrap()
.desired_peers = vec![serde_json::json!({
"deviceId":"space-peer", "online":false,
})];
}
let devices = ["local", "known-offline"]
.into_iter()
.map(|id| {
serde_json::from_value::<openrtc::signaling::Device>(serde_json::json!({
"deviceId":id, "deviceName":"Fixture", "online":false,
}))
.unwrap()
})
.collect::<Vec<_>>();
assert!(submit_native_device_snapshot(
&client,
&state.capability_registry,
&old_source,
"devices",
"local",
&devices,
true
)
.await
.unwrap());
{
let mut registry = state.capability_registry.lock().await;
let peers = ®istry.registrations["devices"].desired_peers;
assert_eq!(peers.len(), 1);
assert_eq!(peers[0]["deviceId"], "known-offline");
assert_eq!(
peers[0]["online"], false,
"offline inventory remains a typed actor input"
);
let (_, aggregate) = registry.aggregate_desired_peers().unwrap();
let aggregate: Vec<serde_json::Value> = serde_json::from_str(&aggregate).unwrap();
assert_eq!(
aggregate.len(),
2,
"native user roster must preserve another avenue"
);
}
let new_source = state
.configure_native_identity("main", "second-user".into(), |_| Ok(()))
.await
.unwrap();
{
let mut registry = state.capability_registry.lock().await;
assert!(matches!(
registry.registrations["devices"].native_source,
Some(NativeRosterSource::Retired)
));
assert!(registry.registrations["devices"].desired_peers.is_empty());
registry
.registrations
.get_mut("devices")
.unwrap()
.native_source = Some(NativeRosterSource::Active(new_source.clone()));
}
assert!(!submit_native_device_snapshot(
&client,
&state.capability_registry,
&old_source,
"devices",
"local",
&[],
true
)
.await
.unwrap());
assert!(submit_native_device_snapshot(
&client,
&state.capability_registry,
&new_source,
"devices",
"local",
&devices,
false
)
.await
.unwrap());
{
let registry = state.capability_registry.lock().await;
assert!(registry.registrations["devices"].desired_peers.is_empty());
assert_eq!(registry.registrations["space"].desired_peers.len(), 1);
}
state
.clear_native_identity("main", &new_source)
.await
.unwrap();
assert!(!submit_native_device_snapshot(
&client,
&state.capability_registry,
&new_source,
"devices",
"local",
&devices,
true
)
.await
.unwrap());
client.stop_external_auto_connect().await;
}
#[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_installed_host_facts() {
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,
if cfg!(feature = "native-broadcast-moq") {
CapabilityMaturity::Preview
} else {
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");
}
#[test]
fn native_testing_endpoint_pair_is_owned_by_the_plugin_state() {
let state = OpenRtcTauriState::new(test_config())
.with_testing_endpoints("http://127.0.0.1:5004", "http://127.0.0.1:8787");
assert_eq!(
state.testing_endpoints,
Some((
"http://127.0.0.1:5004".to_string(),
"http://127.0.0.1:8787".to_string(),
)),
);
}
#[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");
}
}