pub use crate::native_key_store::{load_or_create_endpoint_key, FileSigner};
pub use crate::native_coordination_gateway::{
NativeAuthorityAssignment, NativeAuthorityInterestAssignment, NativeCapabilities,
NativeCapabilityHandle, NativeCapabilityKind, NativeGatewayGrant, NativeGatewayGrantProvider,
NativeGatewayGrantRequest, NativeServiceErrorObservation,
};
#[cfg(feature = "adaptive-room-sentinel")]
#[doc(hidden)]
#[derive(Debug, Clone, Copy, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AdaptiveRoomSentinelOwnerLimits {
pub max_reconnect_attempts: u8,
pub max_operation_attempts: usize,
pub max_prepare_page_members: usize,
pub max_artifact_chunks: usize,
}
#[cfg(feature = "adaptive-room-sentinel")]
#[doc(hidden)]
pub fn adaptive_room_sentinel_owner_limits() -> AdaptiveRoomSentinelOwnerLimits {
AdaptiveRoomSentinelOwnerLimits {
max_reconnect_attempts: crate::native_coordination_gateway::MAX_RECONNECT_ATTEMPTS,
max_operation_attempts: crate::native_coordination_gateway::MAX_OPERATION_ATTEMPTS,
max_prepare_page_members: crate::managed_group_controller::MAX_PREPARE_PAGE_MEMBERS,
max_artifact_chunks: crate::managed_group_controller::MAX_ARTIFACT_CHUNKS,
}
}
#[cfg(feature = "managed-group-encryption")]
pub struct NativeManagedGroupController {
inner: crate::managed_group_controller::ManagedGroupController,
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeManagedProtectedPayload {
pub payload: String,
pub sealed_state: String,
}
#[cfg(feature = "managed-group-encryption")]
impl NativeManagedGroupController {
pub fn new(
device_id: &str,
wrapping_key: [u8; 32],
sealed_state: Option<&[u8]>,
) -> Result<Self> {
Ok(Self {
inner: crate::managed_group_controller::ManagedGroupController::new(
device_id,
wrapping_key,
sealed_state,
)?,
})
}
pub fn publish_key_package(&self) -> Result<Value> {
Ok(serde_json::to_value(self.inner.publish_key_package()?)?)
}
pub fn handle_prepare_page(&mut self, page: Value) -> Result<Vec<Value>> {
let page = serde_json::from_value(page)?;
self.inner
.handle_prepare_page(page)?
.into_iter()
.map(|action| serde_json::to_value(action).map_err(Into::into))
.collect()
}
pub fn handle_artifact_chunk(&mut self, chunk: Value) -> Result<Vec<Value>> {
let chunk = serde_json::from_value(chunk)?;
self.inner
.handle_artifact_chunk(chunk)?
.into_iter()
.map(|action| serde_json::to_value(action).map_err(Into::into))
.collect()
}
#[allow(clippy::too_many_arguments)]
pub fn seal_payload(
&mut self,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: &str,
channel: &str,
priority: u8,
zone_id: Option<&str>,
payload: &[u8],
) -> Result<NativeManagedProtectedPayload> {
let protected = self.inner.seal_payload(
architecture_epoch,
encryption_epoch,
message_id,
channel,
priority,
zone_id,
payload,
)?;
Ok(NativeManagedProtectedPayload {
payload: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(protected.data),
sealed_state: base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(protected.sealed_state),
})
}
#[allow(clippy::too_many_arguments)]
pub fn open_payload(
&mut self,
architecture_epoch: u64,
encryption_epoch: u64,
message_id: &str,
channel: &str,
priority: u8,
zone_id: Option<&str>,
ciphertext: &[u8],
) -> Result<NativeManagedProtectedPayload> {
let protected = self.inner.open_payload(
architecture_epoch,
encryption_epoch,
message_id,
channel,
priority,
zone_id,
ciphertext,
)?;
Ok(NativeManagedProtectedPayload {
payload: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(protected.data),
sealed_state: base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(protected.sealed_state),
})
}
}
#[cfg(feature = "managed-group-encryption")]
pub use crate::native_coordination_gateway::{
InMemoryNativeManagedGroupStateStore, NativeManagedGroupStateStore, NativeManagedRoomMessage,
NativeManagedRoomPublish,
};
use crate::signaling::{
Device, DeviceCapabilities, DeviceEvent, SessionEvent, SignalingBackend, SignalingSession,
};
use anyhow::{anyhow, bail, Context, Result};
use async_trait::async_trait;
use base64::Engine as _;
use ed25519_dalek::{Signature, Verifier, VerifyingKey};
use futures::stream::BoxStream;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use std::fmt;
use std::sync::{Arc, Mutex as StdMutex};
use std::time::{SystemTime, UNIX_EPOCH};
use tokio::sync::{watch, Mutex};
use uuid::Uuid;
pub const OPENRTC_PRODUCTION_CONTROL_PLANE: &str = "https://api.openrtc.app";
const DEVICE_CERTIFICATE_RENEW_SKEW_MS: u64 = 24 * 60 * 60_000;
const SOURCE_RENEW_SKEW_MS: u64 = 5 * 60_000;
const REQUEST_TIMEOUT_SECONDS: u64 = 15;
#[derive(Debug)]
pub struct ControlPlaneHttpError {
pub(crate) status: u16,
pub(crate) message: String,
pub(crate) reason: Option<String>,
pub service_error: Option<crate::service_errors::ServiceError>,
}
impl ControlPlaneHttpError {
pub(crate) fn is_retryable(&self) -> bool {
if let Some(service) = &self.service_error {
return service.retryable;
}
self.status == 408 || self.status == 425 || self.status == 429 || self.status >= 500
}
}
impl fmt::Display for ControlPlaneHttpError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
formatter,
"OpenRTC control-plane request failed ({}): {}",
self.status, self.message
)
}
}
impl std::error::Error for ControlPlaneHttpError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IdentityAssertion {
pub token: String,
pub provider_id: Option<String>,
}
#[async_trait]
pub trait AssertionProvider: Send + Sync {
async fn assertion(&self, force_refresh: bool) -> Result<IdentityAssertion>;
async fn assertion_for_device(
&self,
force_refresh: bool,
_device_id: &str,
) -> Result<IdentityAssertion> {
self.assertion(force_refresh).await
}
async fn assertion_for_device_recovery(&self, device_id: &str) -> Result<IdentityAssertion> {
self.assertion_for_device(true, device_id).await
}
fn session_key(&self) -> Result<Option<String>> {
Ok(None)
}
fn identity_epoch(&self) -> u64 {
0
}
fn subscribe_identity_epoch(&self) -> Option<watch::Receiver<u64>> {
None
}
}
pub trait DeviceSigner: Send + Sync {
fn public_jwk(&self, app_tag: &str) -> Result<Value>;
fn sign(&self, app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>>;
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct StoredCertificate {
pub app_tag: String,
pub principal_id: String,
pub device_id: String,
pub token: String,
pub expires_at_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_key_hash: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub signing_public_jwk: Option<Value>,
}
#[async_trait]
pub trait CertificateStore: Send + Sync {
async fn load(
&self,
app_tag: &str,
principal_id: &str,
device_id: &str,
) -> Result<Option<StoredCertificate>>;
async fn save(&self, certificate: &StoredCertificate) -> Result<()>;
async fn load_for_session(
&self,
_app_tag: &str,
_session_key_hash: &str,
_device_id: &str,
) -> Result<Option<StoredCertificate>> {
Ok(None)
}
async fn remove(&self, app_tag: &str, principal_id: &str, device_id: &str) -> Result<()>;
}
#[derive(Default)]
pub struct InMemoryCertificate {
certificate: StdMutex<Option<StoredCertificate>>,
}
#[async_trait]
impl CertificateStore for InMemoryCertificate {
async fn load(
&self,
app_tag: &str,
principal_id: &str,
device_id: &str,
) -> Result<Option<StoredCertificate>> {
Ok(self
.certificate
.lock()
.map_err(|_| anyhow!("native device certificate cache is poisoned"))?
.clone()
.filter(|value| {
value.app_tag == app_tag
&& value.principal_id == principal_id
&& value.device_id == device_id
}))
}
async fn save(&self, certificate: &StoredCertificate) -> Result<()> {
*self
.certificate
.lock()
.map_err(|_| anyhow!("native device certificate cache is poisoned"))? =
Some(certificate.clone());
Ok(())
}
async fn load_for_session(
&self,
app_tag: &str,
session_key_hash: &str,
device_id: &str,
) -> Result<Option<StoredCertificate>> {
Ok(self
.certificate
.lock()
.map_err(|_| anyhow!("native device certificate cache is poisoned"))?
.clone()
.filter(|value| {
value.app_tag == app_tag
&& value.device_id == device_id
&& value.session_key_hash.as_deref() == Some(session_key_hash)
}))
}
async fn remove(&self, app_tag: &str, principal_id: &str, device_id: &str) -> Result<()> {
let mut guard = self
.certificate
.lock()
.map_err(|_| anyhow!("native device certificate cache is poisoned"))?;
if guard.as_ref().is_some_and(|value| {
value.app_tag == app_tag
&& value.principal_id == principal_id
&& value.device_id == device_id
}) {
*guard = None;
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AttestationEvidence {
pub kind: String,
pub token: String,
}
#[async_trait]
pub trait AttestationProvider: Send + Sync {
async fn evidence(&self, challenge: &str, api_key: &str) -> Result<AttestationEvidence>;
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Features {
pub iroh_relay: bool,
pub managed_turn: bool,
pub moq: bool,
pub ble: bool,
pub advanced_fanout: bool,
pub durable_membership: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct DeviceProfile {
#[serde(skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub platform: Option<String>,
}
#[derive(Clone, Default)]
pub struct DeviceOptions {
pub max_peers: Option<u32>,
pub features: Features,
pub attestation: Option<Arc<dyn AttestationProvider>>,
pub identity_relay: Option<CredentialRelay>,
pub device_profile: Option<DeviceProfile>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CapabilityOptions {
pub max_peers: Option<u32>,
pub features: Features,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RoomArchitectureMode {
#[default]
Auto,
Mesh,
Sparse,
Managed,
Authority,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum RoomDelivery {
#[default]
Reliable,
LatestState,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum EffectiveRoomArchitecture {
Mesh,
Sparse,
Managed,
Authority,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RoomArchitecturePhase {
Preparing,
Settled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RoomArchitectureReason {
Size,
Traffic,
Latency,
Cost,
Capacity,
Manual,
Recovery,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RoomArchitectureSnapshot {
pub requested: RoomArchitectureMode,
pub effective: EffectiveRoomArchitecture,
pub epoch: u64,
pub phase: RoomArchitecturePhase,
pub reason: RoomArchitectureReason,
pub held_credits_usd: f64,
pub quote_expires_at_ms: u64,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RoomOptions {
pub max_peers: Option<u32>,
pub features: Features,
pub architecture: RoomArchitectureMode,
pub delivery: RoomDelivery,
}
pub struct RoomAuthorityOptions {
pub service_id: String,
pub generation: u64,
pub shard_ids: Vec<String>,
pub max_peers: Option<u32>,
pub features: Features,
pub assignment_signer: Arc<dyn DeviceSigner>,
}
pub struct RoomAuthority {
handle: NativeCapabilityHandle,
app_tag: String,
service_id: String,
generation: u64,
shard_ids: Vec<String>,
assignment_signer: Arc<dyn DeviceSigner>,
}
impl RoomAuthority {
pub fn subscribe_service_errors(
&self,
) -> Result<tokio::sync::broadcast::Receiver<NativeServiceErrorObservation>> {
self.handle.subscribe_service_errors()
}
pub async fn compose_client(
&self,
builder: crate::client::ClientBuilder,
) -> Result<Arc<crate::Client>> {
self.handle.compose_authority_client(builder).await
}
pub fn signaling(&self) -> Arc<dyn SignalingBackend> {
self.handle.signaling()
}
pub fn room_architecture(&self) -> Result<Option<RoomArchitectureSnapshot>> {
self.handle.room_architecture()
}
pub async fn publish_assignment(
&self,
revision: u64,
subject_device_id: impl Into<String>,
shard_id: impl Into<String>,
priority_device_ids: Vec<String>,
relevant_entity_ids: Vec<String>,
ttl_ms: u64,
) -> Result<NativeAuthorityInterestAssignment> {
if revision == 0 || ttl_ms == 0 || ttl_ms > 2 * 60_000 {
bail!("authority assignment revision or TTL is invalid");
}
let subject_device_id = bounded_id("subject_device_id", subject_device_id.into())?;
let shard_id = bounded_id("shard_id", shard_id.into())?;
if !self.shard_ids.contains(&shard_id) {
bail!("authority assignment shard is not owned by this service replica");
}
crate::native_coordination_gateway::validate_native_authority_assignment_targets(
&priority_device_ids,
&relevant_entity_ids,
)?;
let issued_at_ms = now_ms();
let assignment = NativeAuthorityInterestAssignment {
policy_version: "room-authority-assignment-v1".to_string(),
service_id: self.service_id.clone(),
generation: self.generation,
revision,
subject_device_id,
shard_id,
priority_device_ids,
relevant_entity_ids,
issued_at_ms,
expires_at_ms: issued_at_ms.saturating_add(ttl_ms),
};
let payload =
serde_json::to_vec(&assignment).context("encode native authority assignment")?;
let signature = self.assignment_signer.sign(&self.app_tag, &payload)?;
if signature.len() != 64 {
bail!("native authority assignment signer returned an invalid Ed25519 signature");
}
self.handle
.publish_authority_assignment(
assignment.clone(),
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature),
)
.await?;
Ok(assignment)
}
pub async fn close(&self) {
self.handle.close().await;
}
}
impl From<CapabilityOptions> for RoomOptions {
fn from(value: CapabilityOptions) -> Self {
Self {
max_peers: value.max_peers,
features: value.features,
architecture: RoomArchitectureMode::Auto,
delivery: RoomDelivery::Reliable,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AnonymousKind {
Space,
Room,
Ticket,
}
impl AnonymousKind {
fn source_kind(self) -> &'static str {
match self {
Self::Space => "space",
Self::Room => "room",
Self::Ticket => "ticket",
}
}
}
#[derive(Clone, Default)]
pub struct CredentialRelay {
value: Arc<StdMutex<Option<String>>>,
}
impl CredentialRelay {
pub fn provider(&self) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
let state = self.value.clone();
Box::new(move || state.lock().ok().and_then(|value| value.clone()))
}
fn set(&self, value: Option<String>) -> Result<()> {
*self
.value
.lock()
.map_err(|_| anyhow!("native identity credential state is poisoned"))? = value;
Ok(())
}
}
#[derive(Default)]
pub struct SignalingSlot {
backend: tokio::sync::RwLock<Option<Arc<dyn SignalingBackend>>>,
}
impl SignalingSlot {
pub async fn install(&self, backend: Arc<dyn SignalingBackend>) -> Result<()> {
let mut current = self.backend.write().await;
if current.is_some() {
bail!("native OpenRTC capability signaling is already installed");
}
*current = Some(backend);
Ok(())
}
async fn current(&self) -> Result<Arc<dyn SignalingBackend>> {
self.backend
.read()
.await
.clone()
.ok_or_else(|| anyhow!("native OpenRTC capability is not active"))
}
}
#[async_trait]
impl SignalingBackend for SignalingSlot {
async fn update_presence(
&self,
user_id: &str,
local_node_id: &str,
ticket_str: &str,
is_online: bool,
name: &str,
ttl_ms: u64,
metadata: Option<&str>,
) -> Result<()> {
self.current()
.await?
.update_presence(
user_id,
local_node_id,
ticket_str,
is_online,
name,
ttl_ms,
metadata,
)
.await
}
async fn set_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
self.current()
.await?
.set_offline(user_id, local_node_id)
.await
}
async fn update_live_presence(
&self,
user_id: &str,
local_node_id: &str,
ticket_str: &str,
name: &str,
metadata: Option<&str>,
) -> Result<()> {
self.current()
.await?
.update_live_presence(user_id, local_node_id, ticket_str, name, metadata)
.await
}
async fn set_live_presence_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
self.current()
.await?
.set_live_presence_offline(user_id, local_node_id)
.await
}
async fn update_device(
&self,
user_id: &str,
device_id: &str,
device_name: Option<&str>,
capabilities: Option<DeviceCapabilities>,
metadata: Option<&str>,
) -> Result<()> {
self.current()
.await?
.update_device(user_id, device_id, device_name, capabilities, metadata)
.await
}
async fn delete_device(&self, user_id: &str, device_id: &str) -> Result<()> {
self.current()
.await?
.delete_device(user_id, device_id)
.await
}
async fn set_excluded_peers(
&self,
user_id: &str,
local_node_id: &str,
excluded_peers: &[String],
) -> Result<()> {
self.current()
.await?
.set_excluded_peers(user_id, local_node_id, excluded_peers)
.await
}
async fn search_devices(
&self,
user_id: &str,
exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
self.current()
.await?
.search_devices(user_id, exclude_node_id)
.await
}
async fn list_devices(
&self,
user_id: &str,
exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
self.current()
.await?
.list_devices(user_id, exclude_node_id)
.await
}
async fn send_message(
&self,
sender_id: &str,
target_id: &str,
payload: &str,
state: Option<&str>,
reply_payload: Option<&str>,
) -> Result<String> {
self.current()
.await?
.send_message(sender_id, target_id, payload, state, reply_payload)
.await
}
async fn subscribe_devices(
&self,
user_id: &str,
) -> Result<BoxStream<'static, Result<Vec<DeviceEvent>>>> {
self.current().await?.subscribe_devices(user_id).await
}
async fn create_session(&self, session: SignalingSession) -> Result<()> {
self.current().await?.create_session(session).await
}
async fn update_session(&self, session_id: &str, update_data: Value) -> Result<()> {
self.current()
.await?
.update_session(session_id, update_data)
.await
}
async fn subscribe_sessions(
&self,
local_device_id: &str,
) -> Result<BoxStream<'static, Result<Vec<SessionEvent>>>> {
self.current()
.await?
.subscribe_sessions(local_device_id)
.await
}
}
#[derive(Debug, Clone)]
struct CredentialSource {
token: String,
expires_at_ms: u64,
principal_id: String,
relay: RelayAvailability,
}
#[derive(Debug, Clone)]
struct AnonymousCredentialSource {
token: String,
expires_at_ms: u64,
principal_id: String,
device_id: String,
requested_id: String,
avenue_id: String,
relay: RelayAvailability,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RelayAvailability {
pub iroh_relay: bool,
pub managed_turn: bool,
}
fn relay_availability(
iroh_relay: Option<bool>,
managed_turn: bool,
legacy_relay: Option<bool>,
) -> RelayAvailability {
RelayAvailability {
iroh_relay: iroh_relay.or(legacy_relay).unwrap_or(true),
managed_turn,
}
}
fn relay_availability_from_token(token: &str) -> Result<RelayAvailability> {
let claims = decode_claims(token)?;
Ok(relay_availability(
claims.get("irohRelay").and_then(Value::as_bool),
claims
.get("managedTurn")
.and_then(Value::as_bool)
.unwrap_or(false),
claims.get("relay").and_then(Value::as_bool),
))
}
fn available_features(mut requested: Features, relay: RelayAvailability) -> Features {
requested.iroh_relay &= relay.iroh_relay;
requested.managed_turn &= relay.managed_turn;
requested
}
#[derive(Clone)]
pub struct ControlPlane {
inner: Arc<ControlPlaneInner>,
}
struct ControlPlaneInner {
api_key: String,
app_tag: String,
endpoint: String,
gateway_endpoint: String,
http: reqwest::Client,
assertion_provider: Option<Arc<dyn AssertionProvider>>,
signer: Arc<dyn DeviceSigner>,
certificate_store: Arc<dyn CertificateStore>,
#[cfg(feature = "managed-group-encryption")]
managed_group_store:
StdMutex<Arc<dyn crate::native_coordination_gateway::NativeManagedGroupStateStore>>,
auth_epoch: Mutex<AuthEpochState>,
}
#[derive(Debug, Clone)]
struct CachedIdentity {
response: AssertionExchangeResponse,
expires_at_ms: u64,
device_id: String,
}
#[derive(Debug, Default)]
struct AuthEpochState {
epoch: Option<u64>,
identity: Option<CachedIdentity>,
}
pub struct Devices {
pub principal_id: String,
pub handle: NativeCapabilityHandle,
identity_relay: CredentialRelay,
identity_epoch_monitor: Option<tokio::task::JoinHandle<()>>,
relay: RelayAvailability,
}
pub struct AnonymousCapability {
pub principal_id: String,
pub handle: NativeCapabilityHandle,
relay: RelayAvailability,
}
impl Devices {
pub fn relay_availability(&self) -> RelayAvailability {
self.relay
}
pub fn identity_credential_provider(&self) -> Box<dyn Fn() -> Option<String> + Send + Sync> {
self.identity_relay.provider()
}
pub fn signaling(&self) -> Arc<dyn SignalingBackend> {
self.handle.signaling()
}
pub fn media_connection(
&self,
client: Arc<crate::Client>,
peer_id: impl Into<String>,
) -> Result<crate::media::MediaConnection> {
media_connection_for_capability(&self.handle, client, peer_id)
}
pub fn is_closed(&self) -> bool {
self.handle.is_closed()
}
pub async fn close(&self) {
if let Some(monitor) = &self.identity_epoch_monitor {
monitor.abort();
}
self.handle.close().await;
}
}
impl Drop for Devices {
fn drop(&mut self) {
if let Some(monitor) = &self.identity_epoch_monitor {
monitor.abort();
}
}
}
impl AnonymousCapability {
pub async fn compose_client(
&self,
builder: crate::client::ClientBuilder,
) -> Result<Arc<crate::Client>> {
self.handle.compose_authority_client(builder).await
}
pub fn relay_availability(&self) -> RelayAvailability {
self.relay
}
pub fn signaling(&self) -> Arc<dyn SignalingBackend> {
self.handle.signaling()
}
pub fn room_architecture(&self) -> Result<Option<RoomArchitectureSnapshot>> {
self.handle.room_architecture()
}
pub fn subscribe_room_architecture(
&self,
) -> Result<watch::Receiver<Option<RoomArchitectureSnapshot>>> {
self.handle.subscribe_room_architecture()
}
pub fn media_connection(
&self,
client: Arc<crate::Client>,
peer_id: impl Into<String>,
) -> Result<crate::media::MediaConnection> {
media_connection_for_capability(&self.handle, client, peer_id)
}
pub fn is_closed(&self) -> bool {
self.handle.is_closed()
}
pub async fn close(&self) {
self.handle.close().await;
}
}
fn media_connection_for_capability(
handle: &NativeCapabilityHandle,
client: Arc<crate::Client>,
peer_id: impl Into<String>,
) -> Result<crate::media::MediaConnection> {
let kind = match handle.kind() {
NativeCapabilityKind::Devices => "devices",
NativeCapabilityKind::Space => "space",
NativeCapabilityKind::Room => "room",
NativeCapabilityKind::Ticket => "ticket",
};
crate::media::MediaConnection::for_capability(client, peer_id, kind, handle.id())
}
impl ControlPlane {
pub fn new(
api_key: &str,
assertion_provider: Arc<dyn AssertionProvider>,
signer: Arc<dyn DeviceSigner>,
certificate_store: Arc<dyn CertificateStore>,
) -> Result<Self> {
let api_key = crate::validate_api_key(api_key)?.to_string();
let app_tag = crate::app_tag_from_api_key(&api_key);
Ok(Self {
inner: Arc::new(ControlPlaneInner {
api_key,
app_tag,
endpoint: OPENRTC_PRODUCTION_CONTROL_PLANE.to_string(),
gateway_endpoint:
crate::native_coordination_gateway::OPENRTC_PRODUCTION_COORDINATION_GATEWAY
.to_string(),
http: reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(REQUEST_TIMEOUT_SECONDS))
.build()?,
assertion_provider: Some(assertion_provider),
signer,
certificate_store,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: StdMutex::new(Arc::new(
crate::native_coordination_gateway::InMemoryNativeManagedGroupStateStore::default(),
)),
auth_epoch: Mutex::new(AuthEpochState::default()),
}),
})
}
pub fn anonymous(api_key: &str, signer: Arc<dyn DeviceSigner>) -> Result<Self> {
let api_key = crate::validate_api_key(api_key)?.to_string();
let app_tag = crate::app_tag_from_api_key(&api_key);
Ok(Self {
inner: Arc::new(ControlPlaneInner {
api_key,
app_tag,
endpoint: OPENRTC_PRODUCTION_CONTROL_PLANE.to_string(),
gateway_endpoint:
crate::native_coordination_gateway::OPENRTC_PRODUCTION_COORDINATION_GATEWAY
.to_string(),
http: reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(REQUEST_TIMEOUT_SECONDS))
.build()?,
assertion_provider: None,
signer,
certificate_store: Arc::new(InMemoryCertificate::default()),
#[cfg(feature = "managed-group-encryption")]
managed_group_store: StdMutex::new(Arc::new(
crate::native_coordination_gateway::InMemoryNativeManagedGroupStateStore::default(),
)),
auth_epoch: Mutex::new(AuthEpochState::default()),
}),
})
}
pub fn app_tag(&self) -> &str {
&self.inner.app_tag
}
#[cfg(feature = "managed-group-encryption")]
pub fn set_managed_group_state_store(
&self,
store: Arc<dyn crate::native_coordination_gateway::NativeManagedGroupStateStore>,
) -> Result<()> {
*self
.inner
.managed_group_store
.lock()
.map_err(|_| anyhow!("native managed group state store lock is poisoned"))? = store;
Ok(())
}
pub fn with_testing_endpoints(
self,
control_plane: impl Into<String>,
gateway: impl Into<String>,
) -> Result<Self> {
#[cfg(any(test, feature = "testing-endpoints"))]
{
let mut control = self;
let control_plane = validate_endpoint(control_plane.into())?;
let gateway = validate_endpoint(gateway.into())?;
Arc::get_mut(&mut control.inner)
.ok_or_else(|| anyhow!("testing endpoints must be set before cloning the client"))?
.endpoint = control_plane;
Arc::get_mut(&mut control.inner)
.ok_or_else(|| anyhow!("testing endpoints must be set before cloning the client"))?
.gateway_endpoint = gateway;
Ok(control)
}
#[cfg(not(any(test, feature = "testing-endpoints")))]
{
let _ = (control_plane.into(), gateway.into());
Err(anyhow!(
"testing endpoints require the OpenRTC testing-endpoints feature"
))
}
}
pub async fn devices(
&self,
device_id: impl Into<String>,
platform_type: impl Into<String>,
options: DeviceOptions,
) -> Result<Devices> {
let device_id = bounded_id("device_id", device_id.into())?;
let platform_type = bounded_id("platform_type", platform_type.into())?;
let device_profile = normalized_device_profile(options.device_profile, &platform_type)?;
let assertion_provider = self.inner.assertion_provider.as_ref().ok_or_else(|| {
anyhow!("authenticated devices require a native identity assertion provider")
})?;
let mut epoch_receiver = assertion_provider.subscribe_identity_epoch();
let identity_epoch = epoch_receiver
.as_ref()
.map(|receiver| *receiver.borrow())
.unwrap_or_else(|| assertion_provider.identity_epoch());
let source = self
.device_certificate(
&device_id,
device_profile.as_ref(),
options.attestation.as_deref(),
identity_epoch,
)
.await?;
let principal_id = source.principal_id.clone();
let relay = source.relay;
let identity_relay = options.identity_relay.clone().unwrap_or_default();
identity_relay.set(Some(source.token.clone()))?;
let grant_provider: Arc<dyn NativeGatewayGrantProvider> =
Arc::new(ControlPlaneGrantProvider {
control_plane: self.clone(),
device_id: device_id.clone(),
source: Mutex::new(source),
max_peers: options.max_peers,
features: options.features,
device_profile,
attestation: options.attestation,
identity_relay: identity_relay.clone(),
identity_epoch,
});
let capabilities = NativeCapabilities::new(
&self.inner.api_key,
device_id,
platform_type,
grant_provider,
)?;
#[cfg(feature = "managed-group-encryption")]
let capabilities = capabilities.with_managed_group_storage(
self.inner.signer.clone(),
self.inner
.managed_group_store
.lock()
.map_err(|_| anyhow!("native managed group state store lock is poisoned"))?
.clone(),
);
#[cfg(any(test, feature = "testing-endpoints"))]
let capabilities = capabilities.with_testing_endpoint(&self.inner.gateway_endpoint)?;
let handle = capabilities.devices(principal_id.clone())?;
let identity_epoch_monitor = epoch_receiver.take().map(|receiver| {
spawn_identity_epoch_monitor(
receiver,
identity_epoch,
handle.closer(),
identity_relay.clone(),
)
});
Ok(Devices {
handle,
principal_id,
identity_relay,
identity_epoch_monitor,
relay,
})
}
pub async fn join_space(
&self,
id: impl Into<String>,
platform_type: impl Into<String>,
options: CapabilityOptions,
) -> Result<AnonymousCapability> {
self.anonymous_capability(
AnonymousKind::Space,
id.into(),
platform_type.into(),
options,
None,
None,
)
.await
}
pub async fn join_room(
&self,
id: impl Into<String>,
platform_type: impl Into<String>,
options: impl Into<RoomOptions>,
) -> Result<AnonymousCapability> {
let options = options.into();
self.anonymous_capability(
AnonymousKind::Room,
id.into(),
platform_type.into(),
CapabilityOptions {
max_peers: options.max_peers,
features: options.features,
},
Some(options.architecture),
Some(options.delivery),
)
.await
}
pub async fn join_authority_room(
&self,
secret_key: impl Into<String>,
room_id: impl Into<String>,
device_id: impl Into<String>,
platform_type: impl Into<String>,
options: RoomAuthorityOptions,
) -> Result<RoomAuthority> {
if options
.max_peers
.is_some_and(|count| !(1..=50).contains(&count))
{
bail!("authority service maxPeers must be within 1..50 room members");
}
let secret_key = required("secret_key", secret_key.into())?;
if !secret_key.starts_with("sk_") || secret_key.len() > 128 {
bail!("authority service secret key is invalid");
}
let room_id = bounded_id("room_id", room_id.into())?;
let device_id = bounded_id("device_id", device_id.into())?;
let platform_type = bounded_id("platform_type", platform_type.into())?;
let service_id = bounded_id("service_id", options.service_id)?;
if service_id.len() > 80 || options.generation == 0 {
bail!("authority service identity or generation is invalid");
}
if options.shard_ids.is_empty() || options.shard_ids.len() > 64 {
bail!("authority service must own from 1 to 64 shards");
}
let mut shard_ids = options
.shard_ids
.into_iter()
.map(|value| bounded_id("shard_id", value))
.collect::<Result<Vec<_>>>()?;
let shard_count = shard_ids.len();
if shard_ids.iter().any(|value| value.len() > 80) {
bail!("authority service shard identifier is invalid");
}
shard_ids.sort();
shard_ids.dedup();
if shard_ids.len() != shard_count {
bail!("authority service shard identifiers must be unique");
}
let assignment_public_jwk = options.assignment_signer.public_jwk(&self.inner.app_tag)?;
validate_public_jwk(&assignment_public_jwk)?;
let grant_provider: Arc<dyn NativeGatewayGrantProvider> =
Arc::new(AuthorityServiceGrantProvider {
control_plane: self.clone(),
secret_key,
service_id: service_id.clone(),
generation: options.generation,
shard_ids: shard_ids.clone(),
assignment_public_jwk,
device_id: device_id.clone(),
max_peers: options.max_peers,
features: options.features,
});
let capabilities = NativeCapabilities::new(
&self.inner.api_key,
device_id,
platform_type,
grant_provider,
)?;
#[cfg(any(test, feature = "testing-endpoints"))]
let capabilities = capabilities.with_testing_endpoint(&self.inner.gateway_endpoint)?;
let handle = capabilities.join_room_with_options(
room_id,
RoomArchitectureMode::Authority,
RoomDelivery::Reliable,
)?;
Ok(RoomAuthority {
handle,
app_tag: self.inner.app_tag.clone(),
service_id,
generation: options.generation,
shard_ids,
assignment_signer: options.assignment_signer,
})
}
pub async fn issue_ticket(
&self,
id: impl Into<String>,
platform_type: impl Into<String>,
options: CapabilityOptions,
) -> Result<AnonymousCapability> {
self.anonymous_capability(
AnonymousKind::Ticket,
id.into(),
platform_type.into(),
options,
None,
None,
)
.await
}
async fn anonymous_capability(
&self,
kind: AnonymousKind,
id: String,
platform_type: String,
options: CapabilityOptions,
architecture: Option<RoomArchitectureMode>,
room_delivery: Option<RoomDelivery>,
) -> Result<AnonymousCapability> {
let platform_type = bounded_id("platform_type", platform_type)?;
let requested_id = bounded_id("capability_id", id)?;
let source = self
.anonymous_source(
kind,
requested_id,
options.max_peers,
architecture,
room_delivery,
)
.await?;
let principal_id = source.principal_id.clone();
let device_id = source.device_id.clone();
let avenue_id = source.avenue_id.clone();
let relay = source.relay;
let grant_provider: Arc<dyn NativeGatewayGrantProvider> =
Arc::new(AnonymousControlPlaneGrantProvider {
control_plane: self.clone(),
source: Mutex::new(source),
kind,
max_peers: options.max_peers,
features: options.features,
architecture,
room_delivery,
});
let capabilities = NativeCapabilities::new(
&self.inner.api_key,
device_id,
platform_type,
grant_provider,
)?;
#[cfg(feature = "managed-group-encryption")]
let capabilities = capabilities.with_managed_group_storage(
self.inner.signer.clone(),
self.inner
.managed_group_store
.lock()
.map_err(|_| anyhow!("native managed group state store lock is poisoned"))?
.clone(),
);
#[cfg(any(test, feature = "testing-endpoints"))]
let capabilities = capabilities.with_testing_endpoint(&self.inner.gateway_endpoint)?;
let handle = match kind {
AnonymousKind::Space => capabilities.join_space(avenue_id)?,
AnonymousKind::Room => capabilities.join_room_with_options(
avenue_id,
architecture.unwrap_or_default(),
room_delivery.unwrap_or_default(),
)?,
AnonymousKind::Ticket => capabilities.issue_ticket(avenue_id)?,
};
Ok(AnonymousCapability {
principal_id,
handle,
relay,
})
}
async fn anonymous_source(
&self,
kind: AnonymousKind,
requested_id: String,
max_peers: Option<u32>,
architecture: Option<RoomArchitectureMode>,
room_delivery: Option<RoomDelivery>,
) -> Result<AnonymousCredentialSource> {
let original_requested_id = requested_id.clone();
let avenue_id = native_capability_avenue_id(&self.inner.api_key, kind, &requested_id);
let proof = self.device_proof(|nonce, issued_at| {
format!(
"openrtc:v2:capability:{}:{}:{}:{}:{}",
self.inner.api_key,
kind.source_kind(),
avenue_id,
nonce,
issued_at
)
})?;
let mut body = json!({
"avenue": { "kind": kind.source_kind(), "id": avenue_id },
"deviceProof": proof,
});
if let Some(max_peers) = max_peers {
body["maxPeers"] = json!(max_peers);
}
if let Some(architecture) = architecture {
body["architecture"] = json!(architecture);
}
if let Some(room_delivery) = room_delivery {
body["roomDelivery"] = json!(room_delivery);
}
let response: CapabilityResponse = self.post("/v2/capabilities", body).await?;
let token = required("capability", response.capability)?;
let principal_id = required_claim_string(&token, "principalId")?;
let thumbprint = required_claim_string(&token, "deviceKeyThumbprint")?;
let claim_exp = decode_claims(&token)?
.get("exp")
.and_then(Value::as_u64)
.ok_or_else(|| anyhow!("OpenRTC capability is missing exp"))?;
if claim_exp != response.expires_at
|| response.avenue.kind != kind.source_kind()
|| response.avenue.id != avenue_id
{
bail!("OpenRTC capability response scope is inconsistent");
}
Ok(AnonymousCredentialSource {
token,
expires_at_ms: response.expires_at.saturating_mul(1_000),
principal_id,
device_id: format!("anon_{}", thumbprint.chars().take(32).collect::<String>()),
requested_id: original_requested_id,
avenue_id,
relay: relay_availability(response.iroh_relay, response.managed_turn, response.relay),
})
}
async fn device_certificate(
&self,
device_id: &str,
device_profile: Option<&DeviceProfile>,
attestation: Option<&dyn AttestationProvider>,
identity_epoch: u64,
) -> Result<CredentialSource> {
let device_key_thumbprint =
native_device_key_thumbprint(&self.inner.signer.public_jwk(&self.inner.app_tag)?)?;
let assertion_provider = self.inner.assertion_provider.as_ref().ok_or_else(|| {
anyhow!("authenticated devices require a native identity assertion provider")
})?;
let session_key_hash = assertion_provider
.session_key()?
.filter(|value| !value.trim().is_empty())
.map(|value| {
native_session_key_hash(&self.inner.app_tag, &device_key_thumbprint, value.trim())
});
if let Some(session_key_hash) = session_key_hash.as_deref() {
if let Some(cached) = self
.inner
.certificate_store
.load_for_session(&self.inner.app_tag, session_key_hash, device_id)
.await?
.filter(|value| certificate_is_reusable(value, &device_key_thumbprint, now_ms()))
{
return Ok(CredentialSource {
relay: relay_availability_from_token(&cached.token)?,
token: cached.token,
expires_at_ms: cached.expires_at_ms,
principal_id: cached.principal_id,
});
}
}
let identity = self
.identity_for_epoch(identity_epoch, device_id, false)
.await?;
let principal_id = bounded_id("principal_id", identity.response.principal_id.clone())?;
if let Some(cached) = self
.inner
.certificate_store
.load(&self.inner.app_tag, &principal_id, device_id)
.await?
.filter(|value| certificate_is_reusable(value, &device_key_thumbprint, now_ms()))
{
return Ok(CredentialSource {
relay: relay_availability_from_token(&cached.token)?,
token: cached.token,
expires_at_ms: cached.expires_at_ms,
principal_id,
});
}
let identity = if identity.expires_at_ms > now_ms().saturating_add(30_000) {
identity
} else {
self.identity_for_epoch(identity_epoch, device_id, true)
.await?
};
if identity.response.principal_id != principal_id {
bail!("native identity principal changed without advancing its login epoch");
}
let enrollment = match self
.enroll_device(
&identity.response.identity_session,
&principal_id,
device_id,
device_profile,
attestation,
false,
)
.await
{
Ok(enrollment) => enrollment,
Err(error) if is_native_device_key_recovery_required(&error) => {
let recovery_identity = self.exchange_recovery_identity(device_id).await?;
if recovery_identity.principal_id != principal_id {
bail!("native identity principal changed during device-key recovery");
}
self.enroll_device(
&recovery_identity.identity_session,
&principal_id,
device_id,
device_profile,
attestation,
true,
)
.await
.context("recover native OpenRTC device key")?
}
Err(error) => return Err(error).context("enroll native OpenRTC device"),
};
let expires_at_ms = enrollment
.expires_at
.checked_mul(1_000)
.ok_or_else(|| anyhow!("device certificate expiry overflow"))?;
let stored = StoredCertificate {
app_tag: self.inner.app_tag.clone(),
principal_id: principal_id.clone(),
device_id: device_id.to_string(),
token: required("device certificate", enrollment.device_certificate)?,
expires_at_ms,
session_key_hash,
signing_public_jwk: Some(enrollment.signing_public_jwk),
};
validate_stored_certificate(&stored, &device_key_thumbprint)?;
self.inner.certificate_store.save(&stored).await?;
Ok(CredentialSource {
relay: relay_availability_from_token(&stored.token)?,
token: stored.token,
expires_at_ms,
principal_id,
})
}
async fn enroll_device(
&self,
identity_session: &str,
principal_id: &str,
device_id: &str,
device_profile: Option<&DeviceProfile>,
attestation: Option<&dyn AttestationProvider>,
recover_device_key: bool,
) -> Result<DeviceEnrollmentResponse> {
let proof = self.device_proof(|nonce, issued_at| {
format!(
"openrtc:v2:device-enroll:{}:{}:{}:{}:{}",
self.inner.app_tag, principal_id, device_id, nonce, issued_at
)
})?;
let attestation_value = if let Some(provider) = attestation {
let challenge = format!(
"openrtc:v2:attestation:{}:{}:{}:{}",
self.inner.app_tag, principal_id, device_id, proof.nonce
);
Some(provider.evidence(&challenge, &self.inner.api_key).await?)
} else {
None
};
let mut body = json!({
"identitySession": identity_session,
"deviceId": device_id,
"deviceProof": proof,
});
if let Some(device_profile) = device_profile {
body["deviceProfile"] = json!(device_profile);
}
if let Some(attestation) = attestation_value {
body["attestation"] = json!(attestation);
}
if recover_device_key {
body["recoverDeviceKey"] = json!(true);
}
self.post("/v2/devices/enroll", body).await
}
async fn exchange_recovery_identity(
&self,
device_id: &str,
) -> Result<AssertionExchangeResponse> {
let assertion_provider = self.inner.assertion_provider.as_ref().ok_or_else(|| {
anyhow!("authenticated devices require a native identity assertion provider")
})?;
let assertion = assertion_provider
.assertion_for_device_recovery(device_id)
.await
.context("obtain consumer device-key recovery assertion")?;
let assertion_token = required("identity assertion", assertion.token)?;
let mut body = json!({ "assertion": assertion_token });
if let Some(provider_id) = assertion.provider_id {
body["providerId"] = json!(provider_id);
}
self.post("/v2/assertions/exchange", body)
.await
.context("exchange consumer device-key recovery assertion")
}
async fn exchange_identity_after_rejection(
&self,
force_assertion_refresh: bool,
device_id: &str,
) -> Result<AssertionExchangeResponse> {
Ok(
match self
.exchange_identity(force_assertion_refresh, device_id)
.await
{
Ok(identity) => identity,
Err(error)
if !force_assertion_refresh
&& error
.downcast_ref::<ControlPlaneHttpError>()
.is_some_and(|failure| failure.status == 401) =>
{
self.exchange_identity(true, device_id).await?
}
Err(error) => return Err(error),
},
)
}
async fn identity_for_epoch(
&self,
identity_epoch: u64,
device_id: &str,
require_unexpired_session: bool,
) -> Result<CachedIdentity> {
let assertion_provider = self.inner.assertion_provider.as_ref().ok_or_else(|| {
anyhow!("authenticated devices require a native identity assertion provider")
})?;
if assertion_provider.identity_epoch() != identity_epoch {
bail!("native identity changed before OpenRTC capability activation");
}
let mut state = self.inner.auth_epoch.lock().await;
if state.epoch != Some(identity_epoch) {
state.epoch = Some(identity_epoch);
state.identity = None;
}
if let Some(identity) = state
.identity
.clone()
.filter(|identity| identity.device_id == device_id)
{
if !require_unexpired_session
|| identity.expires_at_ms > now_ms().saturating_add(30_000)
{
if assertion_provider.identity_epoch() != identity_epoch {
state.identity = None;
bail!("native identity changed during capability activation");
}
return Ok(identity);
}
}
let response = self
.exchange_identity_after_rejection(false, device_id)
.await?;
if assertion_provider.identity_epoch() != identity_epoch {
state.identity = None;
bail!("native identity changed during assertion exchange");
}
let expires_at_ms = response
.expires_at
.checked_mul(1_000)
.ok_or_else(|| anyhow!("identity session expiry overflow"))?;
let identity = CachedIdentity {
response,
expires_at_ms,
device_id: device_id.to_string(),
};
state.identity = Some(identity.clone());
Ok(identity)
}
async fn exchange_identity(
&self,
force_refresh: bool,
device_id: &str,
) -> Result<AssertionExchangeResponse> {
let assertion_provider = self.inner.assertion_provider.as_ref().ok_or_else(|| {
anyhow!("authenticated devices require a native identity assertion provider")
})?;
let assertion = assertion_provider
.assertion_for_device(force_refresh, device_id)
.await
.context("obtain consumer identity assertion")?;
let assertion_token = required("identity assertion", assertion.token)?;
let mut assertion_body = json!({ "assertion": assertion_token });
if let Some(provider_id) = assertion.provider_id {
assertion_body["providerId"] = json!(provider_id);
}
self.post("/v2/assertions/exchange", assertion_body)
.await
.context("exchange consumer identity assertion")
}
fn device_proof(
&self,
challenge: impl FnOnce(&str, u64) -> String,
) -> Result<NativeDeviceProof> {
let public_key_jwk = self.inner.signer.public_jwk(&self.inner.app_tag)?;
validate_public_jwk(&public_key_jwk)?;
let nonce = Uuid::new_v4().simple().to_string();
let issued_at = now_seconds();
let challenge = challenge(&nonce, issued_at);
let signature = self
.inner
.signer
.sign(&self.inner.app_tag, challenge.as_bytes())?;
if signature.len() != 64 {
bail!("native device signer returned an invalid Ed25519 signature");
}
Ok(NativeDeviceProof {
public_key_jwk,
signature: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature),
nonce,
issued_at,
})
}
async fn post<T: for<'de> Deserialize<'de>>(&self, path: &str, body: Value) -> Result<T> {
let mut object = body
.as_object()
.cloned()
.ok_or_else(|| anyhow!("native control-plane body must be an object"))?;
object.insert("apiKey".to_string(), json!(self.inner.api_key));
self.post_json(path, Value::Object(object), None).await
}
async fn post_with_bearer<T: for<'de> Deserialize<'de>>(
&self,
path: &str,
body: Value,
bearer: &str,
) -> Result<T> {
if !bearer.starts_with("sk_") || bearer.len() > 128 {
bail!("native authority service secret is invalid");
}
self.post_json(path, body, Some(bearer)).await
}
async fn post_json<T: for<'de> Deserialize<'de>>(
&self,
path: &str,
body: Value,
bearer: Option<&str>,
) -> Result<T> {
let url = format!("{}{}", self.inner.endpoint.trim_end_matches('/'), path);
let request_id = format!("native_{}", Uuid::new_v4().simple());
let serialized = serde_json::to_vec(&body).context("encode OpenRTC request")?;
for attempt in 0..2 {
let mut request = self
.inner
.http
.post(&url)
.header("X-OpenRTC-Idempotency-Key", &request_id)
.header(reqwest::header::CONTENT_TYPE, "application/json")
.body(serialized.clone());
if let Some(bearer) = bearer {
request = request.bearer_auth(bearer);
}
let response = request.send().await;
let response = match response {
Ok(response) => response,
Err(_error) if attempt == 0 => {
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
continue;
}
Err(error) => return Err(error).context("send OpenRTC request"),
};
let status = response.status();
if !status.is_success() {
let payload = response.json::<Value>().await.unwrap_or(Value::Null);
let message = payload
.get("error")
.and_then(Value::as_str)
.unwrap_or("OpenRTC control-plane request failed")
.chars()
.take(256)
.collect::<String>();
let error = ControlPlaneHttpError {
service_error: crate::service_errors::ServiceError::from_value(&payload),
status: status.as_u16(),
message,
reason: payload
.get("reason")
.and_then(Value::as_str)
.filter(|value| {
!value.is_empty()
&& value.len() <= 80
&& value.bytes().all(|byte| {
byte.is_ascii_lowercase()
|| byte.is_ascii_digit()
|| byte == b'-'
})
})
.map(str::to_string),
};
if attempt == 0
&& error.service_error.is_none()
&& matches!(
status,
reqwest::StatusCode::BAD_GATEWAY
| reqwest::StatusCode::SERVICE_UNAVAILABLE
| reqwest::StatusCode::GATEWAY_TIMEOUT
)
{
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
continue;
}
return Err(anyhow::Error::new(error));
}
return response
.json::<T>()
.await
.context("decode OpenRTC response");
}
unreachable!("native control-plane retry loop has two terminal attempts")
}
}
fn is_native_device_key_recovery_required(error: &anyhow::Error) -> bool {
error
.downcast_ref::<ControlPlaneHttpError>()
.is_some_and(|failure| {
failure.status == 412
&& failure.reason.as_deref() == Some("device-key-recovery-required")
})
}
fn spawn_identity_epoch_monitor(
mut receiver: watch::Receiver<u64>,
identity_epoch: u64,
closer: crate::native_coordination_gateway::NativeCapabilityCloser,
relay: CredentialRelay,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
loop {
if receiver.changed().await.is_err() {
return;
}
if *receiver.borrow_and_update() != identity_epoch {
let _ = relay.set(None);
closer.close().await;
return;
}
}
})
}
struct ControlPlaneGrantProvider {
control_plane: ControlPlane,
device_id: String,
source: Mutex<CredentialSource>,
max_peers: Option<u32>,
features: Features,
device_profile: Option<DeviceProfile>,
attestation: Option<Arc<dyn AttestationProvider>>,
identity_relay: CredentialRelay,
identity_epoch: u64,
}
struct AnonymousControlPlaneGrantProvider {
control_plane: ControlPlane,
source: Mutex<AnonymousCredentialSource>,
kind: AnonymousKind,
max_peers: Option<u32>,
features: Features,
architecture: Option<RoomArchitectureMode>,
room_delivery: Option<RoomDelivery>,
}
struct AuthorityServiceGrantProvider {
control_plane: ControlPlane,
secret_key: String,
service_id: String,
generation: u64,
shard_ids: Vec<String>,
assignment_public_jwk: Value,
device_id: String,
max_peers: Option<u32>,
features: Features,
}
#[async_trait]
impl NativeGatewayGrantProvider for ControlPlaneGrantProvider {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
if request.device_id != self.device_id {
bail!("gateway grant request device does not match enrolled device");
}
let mut source = self.source.lock().await;
if source.expires_at_ms <= now_ms().saturating_add(SOURCE_RENEW_SKEW_MS) {
*source = self
.control_plane
.device_certificate(
&self.device_id,
self.device_profile.as_ref(),
self.attestation.as_deref(),
self.identity_epoch,
)
.await?;
self.identity_relay.set(Some(source.token.clone()))?;
}
let source_jti = required_claim_string(&source.token, "jti")?;
let proof = self.control_plane.device_proof(|nonce, issued_at| {
format!(
"openrtc:v2:gateway-grant:{}:{}:{}:{}:{}:{}",
source_jti,
request.avenue.kind,
request.avenue.id,
request.runtime_instance_id,
nonce,
issued_at
)
})?;
let mut grant_body = json!({
"credentialType": "device-certificate",
"credential": source.token,
"avenue": request.avenue,
"runtimeInstanceId": request.runtime_instance_id,
"ticketFingerprint": request.ticket_fingerprint,
"deviceProof": proof,
"features": available_features(self.features.clone(), source.relay),
});
if let Some(refresh_grant) = request.refresh_grant {
grant_body["refreshGrant"] = json!(refresh_grant);
}
if let Some(max_peers) = self.max_peers {
grant_body["maxPeers"] = json!(max_peers);
}
if let Some(architecture) = request.architecture {
grant_body["architecture"] = json!(architecture);
}
if let Some(room_delivery) = request.room_delivery {
grant_body["roomDelivery"] = json!(room_delivery);
}
if let Some(device_profile) = &self.device_profile {
grant_body["deviceProfile"] = json!(device_profile);
}
let response: GatewayGrantResponse = self
.control_plane
.post("/v2/gateway/grants", grant_body)
.await?;
if response.gateway_url.trim_end_matches('/')
!= self
.control_plane
.inner
.gateway_endpoint
.trim_end_matches('/')
{
bail!("OpenRTC returned an unexpected native gateway origin");
}
Ok(NativeGatewayGrant {
protocol_version: response.protocol_version,
gateway_url: response.gateway_url,
route_key: response.route_key,
token: response.token,
expires_at_ms: response.expires_at_ms,
})
}
}
#[async_trait]
impl NativeGatewayGrantProvider for AnonymousControlPlaneGrantProvider {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
let mut source = self.source.lock().await;
if request.device_id != source.device_id {
bail!("gateway grant request device does not match the capability install key");
}
if source.expires_at_ms <= now_ms().saturating_add(SOURCE_RENEW_SKEW_MS) {
*source = self
.control_plane
.anonymous_source(
self.kind,
source.requested_id.clone(),
self.max_peers,
self.architecture,
self.room_delivery,
)
.await?;
}
let expected_kind = match self.kind {
AnonymousKind::Ticket => "session",
_ => self.kind.source_kind(),
};
if request.avenue.kind != expected_kind || request.avenue.id != source.avenue_id {
bail!("gateway grant request avenue does not match the capability");
}
let source_jti = required_claim_string(&source.token, "jti")?;
let proof = self.control_plane.device_proof(|nonce, issued_at| {
format!(
"openrtc:v2:gateway-grant:{}:{}:{}:{}:{}:{}",
source_jti,
request.avenue.kind,
request.avenue.id,
request.runtime_instance_id,
nonce,
issued_at
)
})?;
let mut grant_body = json!({
"credentialType": "capability",
"credential": source.token,
"avenue": request.avenue,
"runtimeInstanceId": request.runtime_instance_id,
"ticketFingerprint": request.ticket_fingerprint,
"deviceProof": proof,
"features": available_features(self.features.clone(), source.relay),
});
if let Some(refresh_grant) = request.refresh_grant {
grant_body["refreshGrant"] = json!(refresh_grant);
}
if let Some(max_peers) = self.max_peers {
grant_body["maxPeers"] = json!(max_peers);
}
if let Some(architecture) = request.architecture {
grant_body["architecture"] = json!(architecture);
}
if let Some(room_delivery) = request.room_delivery {
grant_body["roomDelivery"] = json!(room_delivery);
}
let response: GatewayGrantResponse = self
.control_plane
.post("/v2/gateway/grants", grant_body)
.await?;
if response.gateway_url.trim_end_matches('/')
!= self
.control_plane
.inner
.gateway_endpoint
.trim_end_matches('/')
{
bail!("OpenRTC returned an unexpected native gateway origin");
}
Ok(NativeGatewayGrant {
protocol_version: response.protocol_version,
gateway_url: response.gateway_url,
route_key: response.route_key,
token: response.token,
expires_at_ms: response.expires_at_ms,
})
}
}
#[async_trait]
impl NativeGatewayGrantProvider for AuthorityServiceGrantProvider {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
if request.device_id != self.device_id
|| request.avenue.kind != "room"
|| request.architecture != Some(RoomArchitectureMode::Authority)
{
bail!("authority service grant request scope is invalid");
}
let assignment_key_x = self
.assignment_public_jwk
.get("x")
.and_then(Value::as_str)
.ok_or_else(|| anyhow!("authority assignment public key is invalid"))?;
let proof = self.control_plane.device_proof(|nonce, issued_at| {
format!(
"openrtc:v2:authority-service-grant:{}:{}:{}:{}:{}:{}:{}:{}:{}",
self.control_plane.inner.app_tag,
request.avenue.id,
self.service_id,
self.generation,
self.shard_ids.join(","),
assignment_key_x,
request.runtime_instance_id,
nonce,
issued_at,
)
})?;
let mut body = json!({
"avenue": request.avenue,
"serviceId": self.service_id,
"generation": self.generation,
"shardIds": self.shard_ids,
"assignmentPublicKeyJwk": self.assignment_public_jwk,
"deviceId": self.device_id,
"runtimeInstanceId": request.runtime_instance_id,
"ticketFingerprint": request.ticket_fingerprint,
"deviceProof": proof,
"features": self.features,
});
if let Some(max_peers) = self.max_peers {
body["maxPeers"] = json!(max_peers);
}
if let Some(refresh_grant) = request.refresh_grant {
body["refreshGrant"] = json!(refresh_grant);
}
let response: GatewayGrantResponse = self
.control_plane
.post_with_bearer("/v2/developer/authority/grants", body, &self.secret_key)
.await?;
if response.gateway_url.trim_end_matches('/')
!= self
.control_plane
.inner
.gateway_endpoint
.trim_end_matches('/')
{
bail!("OpenRTC returned an unexpected native gateway origin");
}
Ok(NativeGatewayGrant {
protocol_version: response.protocol_version,
gateway_url: response.gateway_url,
route_key: response.route_key,
token: response.token,
expires_at_ms: response.expires_at_ms,
})
}
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct NativeDeviceProof {
public_key_jwk: Value,
signature: String,
nonce: String,
issued_at: u64,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AssertionExchangeResponse {
identity_session: String,
principal_id: String,
expires_at: u64,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct DeviceEnrollmentResponse {
device_certificate: String,
expires_at: u64,
signing_public_jwk: Value,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct CapabilityResponse {
capability: String,
expires_at: u64,
avenue: CapabilityAvenueResponse,
#[serde(default)]
iroh_relay: Option<bool>,
#[serde(default)]
managed_turn: bool,
#[serde(default)]
relay: Option<bool>,
}
#[derive(Deserialize)]
struct CapabilityAvenueResponse {
kind: String,
id: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayGrantResponse {
protocol_version: u8,
gateway_url: String,
route_key: String,
token: String,
expires_at_ms: u64,
}
fn now_seconds() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
fn now_ms() -> u64 {
now_seconds().saturating_mul(1_000)
}
fn required(field: &str, value: String) -> Result<String> {
let value = value.trim().to_string();
if value.is_empty() {
bail!("{field} is required");
}
Ok(value)
}
fn bounded_id(field: &str, value: String) -> Result<String> {
let value = required(field, value)?;
if value.len() > 160
|| !value.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
})
{
bail!("{field} is invalid");
}
Ok(value)
}
fn bounded_visible_text(field: &str, value: String, max_len: usize) -> Result<String> {
let value = required(field, value)?;
if value.len() > max_len || value.chars().any(char::is_control) {
bail!("{field} is invalid");
}
Ok(value)
}
fn normalized_device_profile(
profile: Option<DeviceProfile>,
platform_type: &str,
) -> Result<Option<DeviceProfile>> {
let profile = profile.unwrap_or_default();
let name = profile
.name
.map(|value| bounded_visible_text("device_profile.name", value, 120))
.transpose()?;
let platform = Some(bounded_visible_text(
"device_profile.platform",
profile
.platform
.unwrap_or_else(|| platform_type.to_string()),
80,
)?);
Ok(Some(DeviceProfile { name, platform }))
}
fn native_capability_avenue_id(api_key: &str, kind: AnonymousKind, requested_id: &str) -> String {
if kind != AnonymousKind::Space {
return requested_id.to_string();
}
let digest =
<sha2::Sha256 as sha2::Digest>::digest(format!("{api_key}:{requested_id}").as_bytes());
hex::encode(digest)
}
#[cfg(any(test, feature = "testing-endpoints"))]
fn validate_endpoint(value: String) -> Result<String> {
let parsed = reqwest::Url::parse(value.trim())?;
if !matches!(parsed.scheme(), "https" | "http")
|| parsed.host_str().is_none()
|| !parsed.username().is_empty()
|| parsed.password().is_some()
|| parsed.query().is_some()
|| parsed.fragment().is_some()
{
bail!("OpenRTC testing endpoint is invalid");
}
Ok(parsed.as_str().trim_end_matches('/').to_string())
}
fn validate_public_jwk(value: &Value) -> Result<()> {
let object = value
.as_object()
.ok_or_else(|| anyhow!("native device public key is invalid"))?;
let x = object.get("x").and_then(Value::as_str).unwrap_or_default();
if object.get("kty").and_then(Value::as_str) != Some("OKP")
|| object.get("crv").and_then(Value::as_str) != Some("Ed25519")
|| object.contains_key("d")
|| x.len() != 43
|| !x
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
{
bail!("native device public key must be a public Ed25519 JWK");
}
Ok(())
}
fn decode_claims(token: &str) -> Result<Value> {
let payload = token
.split('.')
.nth(1)
.ok_or_else(|| anyhow!("OpenRTC credential is malformed"))?;
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(payload)
.context("decode OpenRTC credential claims")?;
serde_json::from_slice(&bytes).context("parse OpenRTC credential claims")
}
fn required_claim_string(token: &str, name: &str) -> Result<String> {
decode_claims(token)?
.get(name)
.and_then(Value::as_str)
.map(str::to_string)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow!("OpenRTC credential is missing {name}"))
}
fn native_device_key_thumbprint(public_key_jwk: &Value) -> Result<String> {
validate_public_jwk(public_key_jwk)?;
let canonical = json!({
"crv": public_key_jwk.get("crv").and_then(Value::as_str),
"kty": public_key_jwk.get("kty").and_then(Value::as_str),
"x": public_key_jwk.get("x").and_then(Value::as_str),
});
let bytes = serde_json::to_vec(&canonical)?;
let digest = <sha2::Sha256 as sha2::Digest>::digest(bytes);
Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(digest))
}
fn native_session_key_hash(
app_tag: &str,
device_key_thumbprint: &str,
session_key: &str,
) -> String {
let digest = <sha2::Sha256 as sha2::Digest>::digest(
format!(
"openrtc:v2:session\0{}\0{}\0{}",
app_tag, device_key_thumbprint, session_key,
)
.as_bytes(),
);
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(digest)
}
fn validate_stored_certificate(
value: &StoredCertificate,
device_key_thumbprint: &str,
) -> Result<()> {
let signing_public_jwk = value
.signing_public_jwk
.as_ref()
.ok_or_else(|| anyhow!("stored OpenRTC device certificate has no verification key"))?;
verify_device_certificate_signature(&value.token, signing_public_jwk)?;
let claims = decode_claims(&value.token)?;
let issued_at = claims.get("iat").and_then(Value::as_u64);
let expires_at = claims.get("exp").and_then(Value::as_u64);
if claims.get("iss").and_then(Value::as_str) != Some("openrtc")
|| claims.get("aud").and_then(Value::as_str) != Some("openrtc:v2:gateway-grant")
|| claims.get("typ").and_then(Value::as_str) != Some("device-certificate")
|| claims.get("appTag").and_then(Value::as_str) != Some(value.app_tag.as_str())
|| claims.get("principalId").and_then(Value::as_str) != Some(value.principal_id.as_str())
|| claims.get("principalKey").and_then(Value::as_str).is_none()
|| claims.get("deviceId").and_then(Value::as_str) != Some(value.device_id.as_str())
|| claims.get("deviceKeyThumbprint").and_then(Value::as_str) != Some(device_key_thumbprint)
|| claims.get("jti").and_then(Value::as_str).is_none()
|| issued_at.is_none()
|| expires_at.and_then(|exp| exp.checked_mul(1_000)) != Some(value.expires_at_ms)
|| issued_at
.zip(expires_at)
.is_none_or(|(iat, exp)| iat >= exp)
{
bail!("stored OpenRTC device certificate scope is invalid");
}
Ok(())
}
fn verify_device_certificate_signature(token: &str, public_jwk: &Value) -> Result<()> {
let mut parts = token.split('.');
let header_segment = parts
.next()
.ok_or_else(|| anyhow!("invalid certificate JWT"))?;
let claims_segment = parts
.next()
.ok_or_else(|| anyhow!("invalid certificate JWT"))?;
let signature_segment = parts
.next()
.ok_or_else(|| anyhow!("invalid certificate JWT"))?;
if parts.next().is_some() {
bail!("invalid certificate JWT");
}
let header: Value = serde_json::from_slice(
&base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(header_segment)
.context("decode certificate header")?,
)?;
if header.get("alg").and_then(Value::as_str) != Some("EdDSA")
|| public_jwk.get("kty").and_then(Value::as_str) != Some("OKP")
|| public_jwk.get("crv").and_then(Value::as_str) != Some("Ed25519")
|| header.get("kid").and_then(Value::as_str)
!= public_jwk.get("kid").and_then(Value::as_str)
{
bail!("unsupported OpenRTC certificate signing key");
}
let public_key = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(
public_jwk
.get("x")
.and_then(Value::as_str)
.ok_or_else(|| anyhow!("certificate signing key is missing x"))?,
)
.context("decode certificate signing key")?;
let public_key: [u8; 32] = public_key
.try_into()
.map_err(|_| anyhow!("certificate signing key has invalid length"))?;
let verifying_key =
VerifyingKey::from_bytes(&public_key).context("parse certificate signing key")?;
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(signature_segment)
.context("decode certificate signature")?;
let signature = Signature::from_slice(&signature).context("parse certificate signature")?;
verifying_key
.verify(
format!("{header_segment}.{claims_segment}").as_bytes(),
&signature,
)
.context("verify OpenRTC device certificate")
}
fn certificate_is_reusable(
value: &StoredCertificate,
device_key_thumbprint: &str,
at_ms: u64,
) -> bool {
value.expires_at_ms > at_ms.saturating_add(DEVICE_CERTIFICATE_RENEW_SKEW_MS)
&& validate_stored_certificate(value, device_key_thumbprint).is_ok()
}
#[cfg(test)]
mod tests {
use super::*;
use ed25519_dalek::{Signer as _, SigningKey};
use std::sync::atomic::{AtomicU64, Ordering};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
struct TestSigner;
struct NeverGrantProvider;
#[async_trait]
impl NativeGatewayGrantProvider for NeverGrantProvider {
async fn grant(&self, _request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
bail!("inactive epoch-retirement test must not request a gateway grant")
}
}
#[derive(Default)]
struct RecordingAssertionProvider {
force_refreshes: StdMutex<Vec<bool>>,
device_ids: StdMutex<Vec<String>>,
identity_epoch: AtomicU64,
}
#[async_trait]
impl AssertionProvider for RecordingAssertionProvider {
async fn assertion(&self, force_refresh: bool) -> Result<IdentityAssertion> {
self.force_refreshes.lock().unwrap().push(force_refresh);
Ok(IdentityAssertion {
token: if force_refresh { "fresh" } else { "stale" }.to_string(),
provider_id: None,
})
}
async fn assertion_for_device(
&self,
force_refresh: bool,
device_id: &str,
) -> Result<IdentityAssertion> {
self.device_ids.lock().unwrap().push(device_id.to_string());
self.assertion(force_refresh).await
}
fn identity_epoch(&self) -> u64 {
self.identity_epoch.load(Ordering::Acquire)
}
}
impl DeviceSigner for TestSigner {
fn public_jwk(&self, _app_tag: &str) -> Result<Value> {
Ok(json!({
"kty": "OKP",
"crv": "Ed25519",
"x": "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
}))
}
fn sign(&self, _app_tag: &str, _challenge: &[u8]) -> Result<Vec<u8>> {
Ok(vec![0; 64])
}
}
fn signed_certificate(
app_tag: &str,
principal_id: &str,
device_id: &str,
expires_at: u64,
) -> (String, Value) {
let signing_key = SigningKey::from_bytes(&[7_u8; 32]);
let kid = "test-signing-key";
let header = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(
serde_json::to_vec(&json!({ "alg": "EdDSA", "typ": "JWT", "kid": kid })).unwrap(),
);
let claims = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(
serde_json::to_vec(&json!({
"iss": "openrtc",
"aud": "openrtc:v2:gateway-grant",
"typ": "device-certificate",
"appTag": app_tag,
"principalId": principal_id,
"principalKey": "principal-key",
"deviceId": device_id,
"deviceKeyThumbprint": "thumbprint",
"iat": expires_at - 60,
"exp": expires_at,
"jti": "test-jti",
}))
.unwrap(),
);
let signing_input = format!("{header}.{claims}");
let signature = signing_key.sign(signing_input.as_bytes());
let token = format!(
"{signing_input}.{}",
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(signature.to_bytes())
);
let public_jwk = json!({
"kty": "OKP",
"crv": "Ed25519",
"alg": "EdDSA",
"kid": kid,
"x": base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(signing_key.verifying_key().to_bytes()),
});
(token, public_jwk)
}
#[test]
fn cached_certificate_is_scope_bound_and_renews_a_day_early() {
let now = now_ms();
let expires_at = now / 1_000 + 2 * 24 * 60 * 60;
let (token, public_jwk) =
signed_certificate("app_0000000000000000", "principal", "device", expires_at);
let value = StoredCertificate {
app_tag: "app_0000000000000000".to_string(),
principal_id: "principal".to_string(),
device_id: "device".to_string(),
token,
expires_at_ms: expires_at * 1_000,
session_key_hash: Some("session-hash".to_string()),
signing_public_jwk: Some(public_jwk),
};
assert!(certificate_is_reusable(&value, "thumbprint", now));
assert!(!certificate_is_reusable(
&value,
"thumbprint",
value.expires_at_ms - DEVICE_CERTIFICATE_RENEW_SKEW_MS + 1,
));
let mut wrong = value.clone();
wrong.device_id = "other-device".to_string();
assert!(!certificate_is_reusable(&wrong, "thumbprint", now));
assert!(!certificate_is_reusable(&value, "rotated-thumbprint", now));
let mut tampered = value.clone();
tampered.principal_id = "attacker".to_string();
assert!(!certificate_is_reusable(&tampered, "thumbprint", now));
let mut missing_key = value.clone();
missing_key.signing_public_jwk = None;
assert!(!certificate_is_reusable(&missing_key, "thumbprint", now));
}
#[test]
fn device_public_jwk_rejects_private_material() {
assert!(validate_public_jwk(&json!({
"kty": "OKP",
"crv": "Ed25519",
"x": "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
}))
.is_ok());
assert!(validate_public_jwk(&json!({
"kty": "OKP",
"crv": "Ed25519",
"x": "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
"d": "secret",
}))
.is_err());
}
#[test]
fn native_device_profile_defaults_platform_and_rejects_control_text() {
assert_eq!(
normalized_device_profile(None, "ios").unwrap(),
Some(DeviceProfile {
name: None,
platform: Some("ios".to_string()),
}),
);
assert_eq!(
normalized_device_profile(
Some(DeviceProfile {
name: Some("Bryant's iPhone".to_string()),
platform: None,
}),
"ios",
)
.unwrap(),
Some(DeviceProfile {
name: Some("Bryant's iPhone".to_string()),
platform: Some("ios".to_string()),
}),
);
assert!(normalized_device_profile(
Some(DeviceProfile {
name: Some("unsafe\nname".to_string()),
platform: Some("ios".to_string()),
}),
"ios",
)
.is_err());
}
#[test]
fn native_relay_capabilities_do_not_infer_managed_turn_from_legacy_relay() {
let legacy = relay_availability(None, false, Some(true));
assert_eq!(
legacy,
RelayAvailability {
iroh_relay: true,
managed_turn: false,
}
);
let requested = Features {
iroh_relay: true,
managed_turn: true,
..Features::default()
};
assert_eq!(
available_features(requested, legacy),
Features {
iroh_relay: true,
managed_turn: false,
..Features::default()
}
);
}
#[test]
fn native_service_error_retryability_overrides_http_status() {
let mut error = ControlPlaneHttpError {
status: 429,
message: "paused".into(),
reason: None,
service_error: None,
};
assert!(error.is_retryable());
error.service_error = crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "credit-exhausted", "retryable": false, "scope": "account", "requestId": "request-1"
}));
assert!(!error.is_retryable());
error.status = 503;
assert!(!error.is_retryable());
error.service_error = crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-rate-limited", "retryable": true, "retryAfterMs": 2300
}));
assert!(error.is_retryable());
assert_eq!(error.service_error.unwrap().retry_after_ms, Some(2300));
}
#[test]
fn native_device_key_recovery_requires_exact_status_and_reason() {
let required = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
status: 412,
message: "recovery required".to_string(),
reason: Some("device-key-recovery-required".to_string()),
});
assert!(is_native_device_key_recovery_required(&required));
for (status, reason) in [
(403, Some("device-key-recovery-required")),
(412, Some("device-certificate-revoked")),
(412, None),
] {
let unrelated = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
status,
message: "unrelated".to_string(),
reason: reason.map(str::to_string),
});
assert!(!is_native_device_key_recovery_required(&unrelated));
}
}
#[tokio::test]
async fn authority_capacity_is_bounded_before_network_or_runtime_start() {
let control =
ControlPlane::anonymous(crate::test_constants::TEST_API_KEY, Arc::new(TestSigner))
.unwrap();
for max_peers in [
None,
Some(0),
Some(1),
Some(8),
Some(50),
Some(51),
Some(5_000),
] {
let result = control
.join_authority_room(
"sk_test_service",
"room",
"service",
"native",
RoomAuthorityOptions {
service_id: "authority".into(),
generation: 1,
shard_ids: vec!["main".into()],
max_peers,
features: Features::default(),
assignment_signer: Arc::new(TestSigner),
},
)
.await;
if max_peers.is_some_and(|count| !(1..=50).contains(&count)) {
assert!(result.err().unwrap().to_string().contains("maxPeers"));
} else {
result.unwrap().close().await;
}
}
}
#[tokio::test]
async fn anonymous_constructor_is_network_idle_and_rejects_authenticated_devices() {
let control =
ControlPlane::anonymous(crate::test_constants::TEST_API_KEY, Arc::new(TestSigner))
.expect("anonymous v2 control plane");
assert_eq!(
control.app_tag(),
crate::app_tag_from_api_key(crate::test_constants::TEST_API_KEY),
);
let error = control
.device_certificate("device", None, None, 0)
.await
.expect_err("anonymous clients do not own consumer authentication");
assert!(error.to_string().contains("identity assertion provider"));
}
#[tokio::test]
async fn control_plane_retry_reuses_the_exact_body_and_idempotency_key() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let observed = Arc::new(Mutex::new(Vec::<(String, Vec<u8>)>::new()));
let server_observed = observed.clone();
let server = tokio::spawn(async move {
for attempt in 0..2 {
let (mut stream, _) = listener.accept().await.unwrap();
let mut bytes = Vec::new();
let mut buffer = [0_u8; 4096];
let header_end = loop {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
if let Some(index) = bytes.windows(4).position(|part| part == b"\r\n\r\n") {
break index + 4;
}
};
let headers = String::from_utf8_lossy(&bytes[..header_end]).to_string();
let content_length = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().unwrap())
})
.unwrap_or(0);
while bytes.len() < header_end + content_length {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
}
let request_id = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("x-openrtc-idempotency-key")
.then(|| value.trim().to_string())
})
.expect("idempotency header");
server_observed.lock().await.push((
request_id,
bytes[header_end..header_end + content_length].to_vec(),
));
let (status, body) = if attempt == 0 {
("503 Service Unavailable", r#"{"error":"retry"}"#)
} else {
("200 OK", r#"{"ok":true}"#)
};
stream
.write_all(
format!(
"HTTP/1.1 {status}\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len(),
)
.as_bytes(),
)
.await
.unwrap();
}
});
let control =
ControlPlane::anonymous(crate::test_constants::TEST_API_KEY, Arc::new(TestSigner))
.unwrap()
.with_testing_endpoints(format!("http://{address}"), "http://127.0.0.1:1")
.unwrap();
let response: Value = control
.post("/retry", json!({ "intent": "same" }))
.await
.unwrap();
assert_eq!(response, json!({ "ok": true }));
server.await.unwrap();
let observed = observed.lock().await;
assert_eq!(observed.len(), 2);
assert_eq!(observed[0], observed[1]);
}
#[tokio::test]
async fn native_identity_refreshes_once_only_after_unauthorized() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
for attempt in 0..2 {
let (mut stream, _) = listener.accept().await.unwrap();
let mut bytes = Vec::new();
let mut buffer = [0_u8; 4096];
let header_end = loop {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
if let Some(index) = bytes.windows(4).position(|part| part == b"\r\n\r\n") {
break index + 4;
}
};
let headers = String::from_utf8_lossy(&bytes[..header_end]);
let content_length = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().unwrap())
})
.unwrap_or(0);
while bytes.len() < header_end + content_length {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
}
let request_body: Value =
serde_json::from_slice(&bytes[header_end..header_end + content_length])
.unwrap();
let expected_assertion = if attempt == 0 { "stale" } else { "fresh" };
assert_eq!(
request_body.get("assertion").and_then(Value::as_str),
Some(expected_assertion),
);
let (status, body) = if attempt == 0 {
("401 Unauthorized", r#"{"error":"assertion rejected"}"#)
} else {
(
"200 OK",
r#"{"identitySession":"identity","principalId":"principal","expiresAt":4102444800}"#,
)
};
stream
.write_all(
format!(
"HTTP/1.1 {status}\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len(),
)
.as_bytes(),
)
.await
.unwrap();
}
});
let provider = Arc::new(RecordingAssertionProvider::default());
let control = ControlPlane::new(
crate::test_constants::TEST_API_KEY,
provider.clone(),
Arc::new(TestSigner),
Arc::new(InMemoryCertificate::default()),
)
.unwrap()
.with_testing_endpoints(format!("http://{address}"), "http://127.0.0.1:1")
.unwrap();
let identity = control
.exchange_identity_after_rejection(false, "device-a")
.await
.unwrap();
assert_eq!(identity.principal_id, "principal");
assert_eq!(*provider.force_refreshes.lock().unwrap(), vec![false, true]);
server.await.unwrap();
}
#[tokio::test]
async fn native_identity_exchange_is_coalesced_per_host_login_epoch_and_device() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
for attempt in 0..2 {
let (mut stream, _) = listener.accept().await.unwrap();
let mut bytes = Vec::new();
let mut buffer = [0_u8; 4096];
let header_end = loop {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
if let Some(index) = bytes.windows(4).position(|part| part == b"\r\n\r\n") {
break index + 4;
}
};
let headers = String::from_utf8_lossy(&bytes[..header_end]);
let content_length = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().unwrap())
})
.unwrap_or(0);
while bytes.len() < header_end + content_length {
let read = stream.read(&mut buffer).await.unwrap();
assert!(read > 0);
bytes.extend_from_slice(&buffer[..read]);
}
let body = format!(
r#"{{"identitySession":"identity-{attempt}","principalId":"principal-{attempt}","expiresAt":4102444800}}"#,
);
stream
.write_all(
format!(
"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len(),
)
.as_bytes(),
)
.await
.unwrap();
}
});
let provider = Arc::new(RecordingAssertionProvider::default());
let control = ControlPlane::new(
crate::test_constants::TEST_API_KEY,
provider.clone(),
Arc::new(TestSigner),
Arc::new(InMemoryCertificate::default()),
)
.unwrap()
.with_testing_endpoints(format!("http://{address}"), "http://127.0.0.1:1")
.unwrap();
let first = control
.identity_for_epoch(0, "device-a", false)
.await
.unwrap();
let reused = control
.identity_for_epoch(0, "device-a", false)
.await
.unwrap();
assert_eq!(first.response.principal_id, "principal-0");
assert_eq!(reused.response.principal_id, "principal-0");
assert_eq!(*provider.force_refreshes.lock().unwrap(), vec![false]);
let next = control
.identity_for_epoch(0, "device-b", false)
.await
.unwrap();
assert_eq!(next.response.principal_id, "principal-1");
assert_eq!(
*provider.force_refreshes.lock().unwrap(),
vec![false, false]
);
assert_eq!(
*provider.device_ids.lock().unwrap(),
vec!["device-a".to_string(), "device-b".to_string()]
);
server.await.unwrap();
}
#[tokio::test]
async fn native_login_epoch_change_retires_the_old_capability_without_network_work() {
let capabilities = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
"device",
"native",
Arc::new(NeverGrantProvider),
)
.unwrap();
let handle = capabilities.devices("principal").unwrap();
let relay = CredentialRelay::default();
relay.set(Some("device-certificate".to_string())).unwrap();
let credential = relay.provider();
let (epoch_sender, epoch_receiver) = watch::channel(7_u64);
let monitor = spawn_identity_epoch_monitor(epoch_receiver, 7, handle.closer(), relay);
epoch_sender.send(8).unwrap();
tokio::time::timeout(std::time::Duration::from_secs(1), monitor)
.await
.expect("epoch retirement must be prompt")
.unwrap();
assert!(handle.is_closed());
assert_eq!(credential(), None);
}
#[test]
fn native_space_namespace_matches_browser_derivation() {
let api_key = crate::test_constants::TEST_API_KEY;
let requested = "portfolio-cursors";
let expected = hex::encode(<sha2::Sha256 as sha2::Digest>::digest(
format!("{api_key}:{requested}").as_bytes(),
));
assert_eq!(
native_capability_avenue_id(api_key, AnonymousKind::Space, requested),
expected,
);
assert_eq!(
native_capability_avenue_id(api_key, AnonymousKind::Room, "match-123"),
"match-123",
);
}
#[tokio::test]
async fn signaling_slot_installs_once_and_fails_closed_before_activation() {
let slot = SignalingSlot::default();
assert!(slot.search_devices("principal", None).await.is_err());
slot.install(Arc::new(crate::signaling::GatewayRequiredSignalingBackend))
.await
.expect("first capability installs");
assert!(slot
.install(Arc::new(crate::signaling::GatewayRequiredSignalingBackend,))
.await
.is_err());
}
}