#[cfg(feature = "managed-group-encryption")]
use crate::managed_group_controller::{
ManagedArtifactChunk, ManagedGroupAction, ManagedGroupController, ManagedPreparePage,
};
#[cfg(feature = "managed-group-encryption")]
use crate::native::DeviceSigner;
use crate::native::{
ControlPlaneHttpError, EffectiveRoomArchitecture, RoomArchitectureMode, RoomArchitecturePhase,
RoomArchitectureReason, RoomArchitectureSnapshot, RoomDelivery,
};
use crate::signaling::{
Device, DeviceCapabilities, DeviceEvent, SessionEvent, SignalingBackend, SignalingEnvelope,
SignalingSession,
};
use anyhow::{anyhow, bail, Context, Result};
use async_trait::async_trait;
use base64::Engine as _;
use futures::{stream::BoxStream, Sink, SinkExt, StreamExt};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::HashMap;
use std::error::Error as StdError;
use std::fmt;
use std::str::FromStr;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tokio::sync::{broadcast, mpsc, oneshot, watch, RwLock};
use tokio_tungstenite::tungstenite::{
client::IntoClientRequest,
http::{HeaderValue, Uri},
Message,
};
use uuid::Uuid;
#[cfg(feature = "managed-group-encryption")]
use zeroize::Zeroize;
const WIRE_PROTOCOL_VERSION: u8 = 1;
const GRANT_PROTOCOL_VERSION: u8 = 2;
const GATEWAY_PROTOCOL: &str = "openrtc.v2";
const GATEWAY_AUTH_PROTOCOL_PREFIX: &str = "openrtc.auth.";
const REQUEST_TIMEOUT: Duration = Duration::from_secs(15);
const ACK_TIMEOUT: Duration = Duration::from_secs(15);
const AUTH_REFRESH_SKEW_MS: u64 = 5 * 60_000;
const SOCKET_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(10 * 60);
#[cfg(feature = "adaptive-room-sentinel")]
pub(crate) const MAX_RECONNECT_ATTEMPTS: u8 = 12;
#[cfg(not(feature = "adaptive-room-sentinel"))]
const MAX_RECONNECT_ATTEMPTS: u8 = 12;
#[cfg(feature = "adaptive-room-sentinel")]
pub(crate) const MAX_OPERATION_ATTEMPTS: usize = 2;
#[cfg(not(feature = "adaptive-room-sentinel"))]
const MAX_OPERATION_ATTEMPTS: usize = 2;
const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30);
const MAX_EXCLUDED_PEERS: usize = 250;
const MAX_AUTHORITY_PRIORITY_PEERS: usize = 8;
const MAX_AUTHORITY_RELEVANT_ENTITIES: usize = 32;
#[cfg(feature = "managed-group-encryption")]
const MAX_MANAGED_BATCH_MESSAGES: usize = 100;
pub const OPENRTC_PRODUCTION_COORDINATION_GATEWAY: &str = "https://gateway.openrtc.app";
#[derive(Debug)]
struct GatewayConnectError {
message: String,
retryable: bool,
}
impl fmt::Display for GatewayConnectError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.message)
}
}
impl StdError for GatewayConnectError {}
fn gateway_connect_error(message: impl Into<String>, retryable: bool) -> anyhow::Error {
anyhow!(GatewayConnectError {
message: message.into(),
retryable,
})
}
fn is_retryable_gateway_error(error: &anyhow::Error) -> bool {
if let Some(classified) = error.downcast_ref::<GatewayConnectError>() {
return classified.retryable;
}
if let Some(control_plane) = error
.chain()
.find_map(|source| source.downcast_ref::<ControlPlaneHttpError>())
{
return control_plane.is_retryable();
}
true
}
fn operation_requires_budget_renewal(code: &str, retryable: bool) -> bool {
retryable && code == "budget-renewal-required"
}
fn terminal_refresh_error(
lease_error: anyhow::Error,
fallback_error: Option<anyhow::Error>,
) -> anyhow::Error {
fallback_error.unwrap_or(lease_error)
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeCoordinationAvenue {
pub kind: String,
pub id: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeGatewayGrantRequest {
pub avenue: NativeCoordinationAvenue,
pub device_id: String,
pub runtime_instance_id: String,
pub ticket_fingerprint: String,
pub purpose: String,
pub architecture: Option<RoomArchitectureMode>,
pub room_delivery: Option<RoomDelivery>,
pub refresh_grant: Option<String>,
}
#[async_trait]
pub trait NativeGatewayGrantProvider: Send + Sync {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant>;
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeGatewayGrant {
pub protocol_version: u8,
pub gateway_url: String,
pub route_key: String,
pub token: String,
pub expires_at_ms: u64,
}
pub struct GatewayOptions {
pub endpoint: String,
pub app_tag: String,
pub device_id: String,
pub platform_type: String,
pub avenue: NativeCoordinationAvenue,
pub architecture: Option<RoomArchitectureMode>,
pub room_delivery: Option<RoomDelivery>,
pub grant_provider: Arc<dyn NativeGatewayGrantProvider>,
#[cfg(feature = "managed-group-encryption")]
pub managed_group_signer: Option<Arc<dyn DeviceSigner>>,
#[cfg(feature = "managed-group-encryption")]
pub managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
}
#[cfg(feature = "managed-group-encryption")]
#[async_trait]
pub trait NativeManagedGroupStateStore: Send + Sync {
async fn load(
&self,
app_tag: &str,
device_id: &str,
avenue_key: &str,
) -> Result<Option<Vec<u8>>>;
async fn save(
&self,
app_tag: &str,
device_id: &str,
avenue_key: &str,
sealed_state: &[u8],
) -> Result<()>;
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Default)]
pub struct InMemoryNativeManagedGroupStateStore {
states: Mutex<HashMap<String, Vec<u8>>>,
}
#[cfg(feature = "managed-group-encryption")]
#[async_trait]
impl NativeManagedGroupStateStore for InMemoryNativeManagedGroupStateStore {
async fn load(
&self,
app_tag: &str,
device_id: &str,
avenue_key: &str,
) -> Result<Option<Vec<u8>>> {
Ok(self
.states
.lock()
.map_err(|_| anyhow!("managed group state store is poisoned"))?
.get(&format!("{app_tag}:{device_id}:{avenue_key}"))
.cloned())
}
async fn save(
&self,
app_tag: &str,
device_id: &str,
avenue_key: &str,
sealed_state: &[u8],
) -> Result<()> {
if sealed_state.is_empty() {
bail!("managed group sealed state is empty");
}
self.states
.lock()
.map_err(|_| anyhow!("managed group state store is poisoned"))?
.insert(
format!("{app_tag}:{device_id}:{avenue_key}"),
sealed_state.to_vec(),
);
Ok(())
}
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeManagedRoomPublish {
pub message_id: String,
pub channel: String,
pub priority: u8,
pub zone_id: Option<String>,
pub payload: Vec<u8>,
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeManagedRoomMessage {
pub sender_device_id: String,
pub message_id: String,
pub channel: String,
pub priority: u8,
pub zone_id: Option<String>,
pub payload: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeAuthorityInterestAssignment {
pub policy_version: String,
pub service_id: String,
pub generation: u64,
pub revision: u64,
pub subject_device_id: String,
pub shard_id: String,
pub priority_device_ids: Vec<String>,
pub relevant_entity_ids: Vec<String>,
pub issued_at_ms: u64,
pub expires_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeAuthorityAssignment {
pub assignment: NativeAuthorityInterestAssignment,
pub signature: String,
}
pub struct NativeCapabilities {
endpoint: String,
app_tag: String,
device_id: String,
platform_type: String,
grant_provider: Arc<dyn NativeGatewayGrantProvider>,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: Option<Arc<dyn DeviceSigner>>,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
}
impl NativeCapabilities {
pub fn new(
api_key: &str,
device_id: impl Into<String>,
platform_type: impl Into<String>,
grant_provider: Arc<dyn NativeGatewayGrantProvider>,
) -> Result<Self> {
let api_key = crate::validate_api_key(api_key)?;
Ok(Self {
endpoint: OPENRTC_PRODUCTION_COORDINATION_GATEWAY.to_string(),
app_tag: crate::app_tag_from_api_key(api_key),
device_id: required("device_id", device_id.into())?,
platform_type: required("platform_type", platform_type.into())?,
grant_provider,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
}
#[cfg(feature = "managed-group-encryption")]
pub fn with_managed_group_storage(
mut self,
signer: Arc<dyn DeviceSigner>,
store: Arc<dyn NativeManagedGroupStateStore>,
) -> Self {
self.managed_group_signer = Some(signer);
self.managed_group_store = Some(store);
self
}
#[cfg(any(test, feature = "testing-endpoints"))]
pub fn with_testing_endpoint(mut self, endpoint: impl Into<String>) -> Result<Self> {
self.endpoint = validate_endpoint("endpoint", endpoint.into(), true)?;
Ok(self)
}
pub fn devices(&self, principal_id: impl Into<String>) -> Result<NativeCapabilityHandle> {
self.open("user", principal_id, None, None)
}
pub fn join_space(&self, id: impl Into<String>) -> Result<NativeCapabilityHandle> {
self.open("space", id, None, None)
}
pub fn join_room_with_architecture(
&self,
id: impl Into<String>,
architecture: RoomArchitectureMode,
) -> Result<NativeCapabilityHandle> {
self.join_room_with_options(id, architecture, RoomDelivery::Reliable)
}
pub fn join_room_with_options(
&self,
id: impl Into<String>,
architecture: RoomArchitectureMode,
delivery: RoomDelivery,
) -> Result<NativeCapabilityHandle> {
self.open("room", id, Some(architecture), Some(delivery))
}
pub fn issue_ticket(&self, id: impl Into<String>) -> Result<NativeCapabilityHandle> {
self.open("session", id, None, None)
}
fn open(
&self,
kind: &'static str,
id: impl Into<String>,
architecture: Option<RoomArchitectureMode>,
room_delivery: Option<RoomDelivery>,
) -> Result<NativeCapabilityHandle> {
NativeCapabilityHandle::new(GatewayOptions {
endpoint: self.endpoint.clone(),
app_tag: self.app_tag.clone(),
device_id: self.device_id.clone(),
platform_type: self.platform_type.clone(),
avenue: NativeCoordinationAvenue {
kind: kind.to_string(),
id: id.into(),
},
architecture,
room_delivery,
grant_provider: self.grant_provider.clone(),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: self.managed_group_signer.clone(),
#[cfg(feature = "managed-group-encryption")]
managed_group_store: self.managed_group_store.clone(),
})
}
}
pub struct NativeCapabilityHandle {
kind: NativeCapabilityKind,
id: String,
signaling: Arc<NativeCoordinationGatewaySignaling>,
closed: Arc<AtomicBool>,
}
pub(crate) struct NativeCapabilityCloser {
signaling: Arc<NativeCoordinationGatewaySignaling>,
closed: Arc<AtomicBool>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NativeCapabilityKind {
Devices,
Space,
Room,
Ticket,
}
impl NativeCapabilityKind {
fn from_avenue_kind(value: &str) -> Result<Self> {
match value {
"user" => Ok(Self::Devices),
"space" => Ok(Self::Space),
"room" => Ok(Self::Room),
"session" => Ok(Self::Ticket),
_ => bail!("native coordination avenue kind is invalid"),
}
}
}
impl NativeCapabilityHandle {
pub fn new(options: GatewayOptions) -> Result<Self> {
let kind = NativeCapabilityKind::from_avenue_kind(&options.avenue.kind)?;
let id = required("avenue.id", options.avenue.id.clone())?;
let signaling = NativeCoordinationGatewaySignaling::new(options)?;
Ok(Self {
kind,
id,
signaling,
closed: Arc::new(AtomicBool::new(false)),
})
}
pub fn kind(&self) -> NativeCapabilityKind {
self.kind
}
pub fn id(&self) -> &str {
&self.id
}
pub fn signaling(&self) -> Arc<dyn SignalingBackend> {
self.signaling.clone()
}
pub fn room_architecture(&self) -> Result<Option<RoomArchitectureSnapshot>> {
if self.kind != NativeCapabilityKind::Room {
bail!("native room architecture is available only for room capabilities");
}
Ok(self.signaling.room_architecture())
}
pub fn subscribe_room_architecture(
&self,
) -> Result<tokio::sync::watch::Receiver<Option<RoomArchitectureSnapshot>>> {
if self.kind != NativeCapabilityKind::Room {
bail!("native room architecture is available only for room capabilities");
}
Ok(self.signaling.subscribe_room_architecture())
}
pub async fn publish_authority_assignment(
&self,
assignment: NativeAuthorityInterestAssignment,
signature: impl Into<String>,
) -> Result<()> {
if self.kind != NativeCapabilityKind::Room {
bail!("authority assignments are available only for room capabilities");
}
let signature = signature.into();
validate_authority_assignment(&assignment, &signature)?;
self.signaling
.publish_authority_assignment(assignment, signature)
.await
}
pub fn authority_assignment(&self) -> Result<Option<NativeAuthorityAssignment>> {
if self.kind != NativeCapabilityKind::Room {
bail!("authority assignments are available only for room capabilities");
}
Ok(self.signaling.authority_assignment())
}
pub fn subscribe_authority_assignments(
&self,
) -> Result<watch::Receiver<Option<NativeAuthorityAssignment>>> {
if self.kind != NativeCapabilityKind::Room {
bail!("authority assignments are available only for room capabilities");
}
Ok(self.signaling.subscribe_authority_assignments())
}
#[cfg(feature = "managed-group-encryption")]
pub async fn publish_managed_room(&self, batch: Vec<NativeManagedRoomPublish>) -> Result<()> {
if self.kind != NativeCapabilityKind::Room {
bail!("managed fanout is available only for room capabilities");
}
self.signaling.publish_managed_room(batch).await
}
#[cfg(feature = "managed-group-encryption")]
pub fn subscribe_managed_room(&self) -> Result<broadcast::Receiver<NativeManagedRoomMessage>> {
if self.kind != NativeCapabilityKind::Room {
bail!("managed fanout is available only for room capabilities");
}
Ok(self.signaling.subscribe_managed_room())
}
pub fn is_closed(&self) -> bool {
self.closed.load(Ordering::Acquire)
}
pub async fn close(&self) {
if !self.closed.swap(true, Ordering::AcqRel) {
self.signaling.stop().await;
}
}
pub(crate) fn closer(&self) -> NativeCapabilityCloser {
NativeCapabilityCloser {
signaling: self.signaling.clone(),
closed: self.closed.clone(),
}
}
}
pub(crate) fn validate_native_authority_assignment_targets(
priority_device_ids: &[String],
relevant_entity_ids: &[String],
) -> Result<()> {
let valid_id = |value: &str, max: usize| {
!value.is_empty()
&& value.len() <= max
&& value.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'.' | b':' | b'@' | b'-')
})
};
if priority_device_ids.len() > MAX_AUTHORITY_PRIORITY_PEERS
|| relevant_entity_ids.len() > MAX_AUTHORITY_RELEVANT_ENTITIES
|| priority_device_ids
.iter()
.any(|value| !valid_id(value, 160))
|| relevant_entity_ids
.iter()
.any(|value| !valid_id(value, 160))
|| priority_device_ids
.iter()
.collect::<std::collections::HashSet<_>>()
.len()
!= priority_device_ids.len()
|| relevant_entity_ids
.iter()
.collect::<std::collections::HashSet<_>>()
.len()
!= relevant_entity_ids.len()
{
bail!("native authority assignment targets are invalid");
}
Ok(())
}
fn validate_authority_assignment(
assignment: &NativeAuthorityInterestAssignment,
signature: &str,
) -> Result<()> {
validate_native_authority_assignment_targets(
&assignment.priority_device_ids,
&assignment.relevant_entity_ids,
)?;
let valid_id = |value: &str, max: usize| {
!value.is_empty()
&& value.len() <= max
&& value.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'.' | b':' | b'@' | b'-')
})
};
if assignment.policy_version != "room-authority-assignment-v1"
|| !valid_id(&assignment.service_id, 80)
|| !valid_id(&assignment.subject_device_id, 160)
|| !valid_id(&assignment.shard_id, 80)
|| assignment.generation == 0
|| assignment.revision == 0
|| assignment.expires_at_ms <= assignment.issued_at_ms
|| assignment.issued_at_ms > now_ms().saturating_add(30_000)
|| assignment.expires_at_ms <= now_ms()
|| assignment
.expires_at_ms
.saturating_sub(assignment.issued_at_ms)
> 2 * 60_000
{
bail!("native authority assignment is invalid");
}
let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(signature)
.context("decode native authority assignment signature")?;
if decoded.len() != 64 {
bail!("native authority assignment signature is invalid");
}
Ok(())
}
impl NativeCapabilityCloser {
pub(crate) async fn close(&self) {
if !self.closed.swap(true, Ordering::AcqRel) {
self.signaling.stop().await;
}
}
}
impl Drop for NativeCapabilityHandle {
fn drop(&mut self) {
if !self.closed.swap(true, Ordering::AcqRel) {
let _ = self.signaling.commands.try_send(Command::Stop);
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct DesiredPresence {
user_id: String,
local_node_id: String,
ticket: String,
device_name: String,
metadata: Option<String>,
ttl_ms: u64,
online: bool,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
struct DevicePatch {
#[serde(skip_serializing_if = "Option::is_none")]
device_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
capabilities: Option<DeviceCapabilities>,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
excluded_peers: Option<Vec<String>>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct OutboundGatewayDevice {
device_id: String,
runtime_instance_id: String,
node_id: String,
device_name: String,
platform_type: String,
ticket: String,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
capabilities: Option<DeviceCapabilities>,
excluded_peers: Vec<String>,
online: bool,
}
#[derive(Debug)]
enum Command {
Publish {
desired: DesiredPresence,
reply: oneshot::Sender<Result<()>>,
},
Patch {
patch: DevicePatch,
reply: oneshot::Sender<Result<()>>,
},
Offline {
reply: oneshot::Sender<Result<()>>,
},
Delete {
user_id: String,
device_id: String,
reply: oneshot::Sender<Result<()>>,
},
SendSignal {
target_device_id: String,
payload: String,
state: Option<String>,
reply_payload: Option<String>,
reply: oneshot::Sender<Result<String>>,
},
PutSession {
session_id: String,
session: serde_json::Value,
expires_at_ms: i64,
reply: oneshot::Sender<Result<()>>,
},
PublishAuthorityAssignment {
assignment: NativeAuthorityInterestAssignment,
signature: String,
reply: oneshot::Sender<Result<()>>,
},
#[cfg(feature = "managed-group-encryption")]
PublishManaged {
batch: Vec<NativeManagedRoomPublish>,
reply: oneshot::Sender<Result<()>>,
},
Stop,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayDevice {
#[serde(default)]
user_id: Option<String>,
device_id: String,
runtime_instance_id: String,
node_id: String,
device_name: String,
platform_type: String,
ticket: String,
#[serde(default)]
metadata: Option<String>,
#[serde(default)]
capabilities: Option<DeviceCapabilities>,
#[serde(default)]
excluded_peers: Vec<String>,
online: bool,
updated_at_ms: i64,
expires_at_ms: i64,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayMember {
#[serde(default)]
user_id: Option<String>,
device_id: String,
device_name: String,
platform_type: String,
#[serde(default)]
metadata: Option<String>,
#[serde(default)]
capabilities: Option<DeviceCapabilities>,
online: bool,
updated_at_ms: i64,
expires_at_ms: i64,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayArchitectureLease {
policy_version: String,
epoch: u64,
#[serde(default)]
previous_architecture_epoch: Option<u64>,
requested_mode: RoomArchitectureMode,
effective_mode: EffectiveRoomArchitecture,
phase: RoomArchitecturePhase,
reason: RoomArchitectureReason,
held_credits_microusd: u64,
quote_expires_at_ms: u64,
#[serde(default)]
reservation_id: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayGroupEncryptionLease {
policy_version: String,
epoch: u64,
#[serde(default)]
previous_encryption_epoch: Option<u64>,
phase: String,
committer_device_id: String,
group_id_hash: String,
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct NativeManagedEncryptedEnvelope {
message_id: String,
channel: String,
priority: u8,
scope: NativeManagedScope,
ciphertext: String,
}
#[cfg(feature = "managed-group-encryption")]
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "lowercase")]
enum NativeManagedScope {
Global,
Zone {
#[serde(rename = "zoneId")]
zone_id: String,
},
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayTopologyLease {
schema_version: u8,
topology_revision: u64,
#[serde(default)]
previous_topology_revision: Option<u64>,
avenue: NativeCoordinationAvenue,
grant_jti: String,
expires_at_ms: u64,
architecture: GatewayArchitectureLease,
#[serde(default)]
group_encryption: Option<GatewayGroupEncryptionLease>,
active: Vec<GatewayDevice>,
backups: Vec<GatewayDevice>,
}
#[derive(Debug, Deserialize)]
#[serde(tag = "type")]
enum ServerFrame {
#[serde(rename = "ready")]
Ready {
#[serde(rename = "budgetRemainingMicrousd")]
_budget_remaining_microusd: i64,
#[serde(default, rename = "leaseRefreshMode")]
lease_refresh_mode: Option<String>,
},
#[serde(rename = "auth.refreshed")]
AuthRefreshed {
#[serde(rename = "expiresAtMs")]
_expires_at_ms: u64,
},
#[serde(rename = "lease.refreshed")]
LeaseRefreshed {
#[serde(rename = "idempotencyKey")]
idempotency_key: String,
#[serde(rename = "expiresAtMs")]
expires_at_ms: u64,
},
#[serde(rename = "roster.snapshot")]
RosterSnapshot { devices: Vec<GatewayDevice> },
#[serde(rename = "membership.snapshot")]
MembershipSnapshot { members: Vec<GatewayMember> },
#[serde(rename = "membership.page")]
MembershipPage {
#[serde(rename = "snapshotId")]
snapshot_id: String,
#[serde(rename = "pageIndex")]
page_index: usize,
#[serde(rename = "pageCount")]
page_count: usize,
#[serde(rename = "memberCount")]
member_count: usize,
members: Vec<GatewayMember>,
},
#[serde(rename = "membership.changed")]
MembershipChanged {
operation: String,
member: GatewayMember,
},
#[serde(rename = "topology.lease")]
TopologyLease { lease: GatewayTopologyLease },
#[cfg(feature = "managed-group-encryption")]
#[serde(rename = "managed.prepare.page")]
ManagedPreparePage {
#[serde(flatten)]
page: ManagedPreparePage,
},
#[cfg(feature = "managed-group-encryption")]
#[serde(rename = "managed.commit.chunk.received")]
ManagedArtifactChunk {
#[serde(flatten)]
chunk: ManagedArtifactChunk,
},
#[cfg(feature = "managed-group-encryption")]
#[serde(rename = "managed.received")]
ManagedReceived {
#[serde(rename = "senderDeviceId")]
sender_device_id: String,
#[serde(rename = "architectureEpoch")]
architecture_epoch: u64,
#[serde(rename = "encryptionEpoch")]
encryption_epoch: u64,
batch: Vec<NativeManagedEncryptedEnvelope>,
},
#[serde(rename = "authority.assignment")]
AuthorityAssignment {
assignment: NativeAuthorityInterestAssignment,
signature: String,
},
#[serde(rename = "presence.changed")]
PresenceChanged {
operation: String,
device: GatewayDevice,
},
#[serde(rename = "session.changed")]
SessionChanged {
operation: String,
#[serde(rename = "sessionId")]
session_id: String,
#[serde(default)]
session: Option<serde_json::Value>,
},
#[serde(rename = "ack")]
Ack {
#[serde(rename = "idempotencyKey")]
idempotency_key: String,
},
#[serde(rename = "error")]
Error {
code: String,
message: String,
#[serde(default, rename = "idempotencyKey")]
idempotency_key: Option<String>,
retryable: bool,
},
#[serde(rename = "signal.received")]
SignalReceived {
#[serde(rename = "signalId")]
_signal_id: String,
#[serde(rename = "senderDeviceId")]
sender_device_id: String,
payload: String,
#[serde(default)]
state: Option<String>,
#[serde(default, rename = "replyPayload")]
reply_payload: Option<String>,
#[serde(rename = "createdAtMs")]
_created_at_ms: i64,
},
#[serde(rename = "pong")]
Pong,
}
struct SharedState {
desired: RwLock<Option<DesiredPresence>>,
applied_presence: RwLock<Option<DesiredPresence>>,
staged_patch: RwLock<DevicePatch>,
members: RwLock<HashMap<String, GatewayMember>>,
membership_pages: Mutex<Option<PendingMembershipPages>>,
topology_routes: RwLock<HashMap<String, GatewayDevice>>,
topology_revision: RwLock<u64>,
encryption_epoch: RwLock<u64>,
active_grant: RwLock<Option<(String, u64)>>,
devices: RwLock<HashMap<String, Device>>,
device_events: broadcast::Sender<Vec<DeviceEvent>>,
session_events: broadcast::Sender<Vec<SessionEvent>>,
room_architecture: watch::Sender<Option<RoomArchitectureSnapshot>>,
authority_assignment: watch::Sender<Option<NativeAuthorityAssignment>>,
pending_messages: Mutex<Vec<SignalingEnvelope>>,
#[cfg(feature = "managed-group-encryption")]
managed_group: tokio::sync::Mutex<Option<ManagedGroupController>>,
#[cfg(feature = "managed-group-encryption")]
managed_messages: broadcast::Sender<NativeManagedRoomMessage>,
}
struct PendingMembershipPages {
snapshot_id: String,
next_page_index: usize,
page_count: usize,
member_count: usize,
members: HashMap<String, GatewayMember>,
}
pub(crate) struct NativeCoordinationGatewaySignaling {
endpoint: String,
app_tag: String,
device_id: String,
avenue: NativeCoordinationAvenue,
architecture: Option<RoomArchitectureMode>,
room_delivery: Option<RoomDelivery>,
runtime_instance_id: String,
platform_type: String,
grant_provider: Arc<dyn NativeGatewayGrantProvider>,
shared: Arc<SharedState>,
commands: mpsc::Sender<Command>,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: Option<Arc<dyn DeviceSigner>>,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
}
impl NativeCoordinationGatewaySignaling {
pub(crate) fn new(options: GatewayOptions) -> Result<Arc<Self>> {
let endpoint = validate_endpoint("endpoint", options.endpoint, true)?;
let app_tag = required("app_tag", options.app_tag)?;
let device_id = required("device_id", options.device_id)?;
let platform_type = required("platform_type", options.platform_type)?;
let avenue = options.avenue;
if !matches!(avenue.kind.as_str(), "user" | "space" | "room" | "session") {
bail!("native coordination avenue kind is invalid");
}
let avenue_id = required("avenue.id", avenue.id)?;
if options.architecture.is_some() && avenue.kind != "room" {
bail!("native room architecture is valid only for room avenues");
}
if options.room_delivery.is_some() && avenue.kind != "room" {
bail!("native room delivery is valid only for room avenues");
}
let (device_events, _) = broadcast::channel(64);
let (session_events, _) = broadcast::channel(64);
let (room_architecture, _) = watch::channel(None);
let (authority_assignment, _) = watch::channel(None);
#[cfg(feature = "managed-group-encryption")]
let (managed_messages, _) = broadcast::channel(256);
let shared = Arc::new(SharedState {
desired: RwLock::new(None),
applied_presence: RwLock::new(None),
staged_patch: RwLock::new(DevicePatch::default()),
members: RwLock::new(HashMap::new()),
membership_pages: Mutex::new(None),
topology_routes: RwLock::new(HashMap::new()),
topology_revision: RwLock::new(0),
encryption_epoch: RwLock::new(0),
active_grant: RwLock::new(None),
devices: RwLock::new(HashMap::new()),
device_events,
session_events,
room_architecture,
authority_assignment,
pending_messages: Mutex::new(Vec::new()),
#[cfg(feature = "managed-group-encryption")]
managed_group: tokio::sync::Mutex::new(None),
#[cfg(feature = "managed-group-encryption")]
managed_messages,
});
let (commands, receiver) = mpsc::channel(64);
let adapter = Arc::new(Self {
endpoint,
app_tag,
device_id,
avenue: NativeCoordinationAvenue {
kind: avenue.kind,
id: avenue_id,
},
architecture: options.architecture,
room_delivery: options.room_delivery,
runtime_instance_id: format!("runtime:{}", Uuid::new_v4()),
platform_type,
grant_provider: options.grant_provider,
shared,
commands,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: options.managed_group_signer,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: options.managed_group_store,
});
tokio::spawn(run_actor(adapter.clone(), receiver));
Ok(adapter)
}
pub async fn stop(&self) {
let _ = self.commands.send(Command::Stop).await;
}
fn room_architecture(&self) -> Option<RoomArchitectureSnapshot> {
self.shared.room_architecture.borrow().clone()
}
fn subscribe_room_architecture(&self) -> watch::Receiver<Option<RoomArchitectureSnapshot>> {
self.shared.room_architecture.subscribe()
}
async fn publish_authority_assignment(
&self,
assignment: NativeAuthorityInterestAssignment,
signature: String,
) -> Result<()> {
validate_authority_assignment(&assignment, &signature)?;
self.request_unit(|reply| Command::PublishAuthorityAssignment {
assignment,
signature,
reply,
})
.await
}
fn authority_assignment(&self) -> Option<NativeAuthorityAssignment> {
self.shared
.authority_assignment
.borrow()
.clone()
.filter(|value| value.assignment.expires_at_ms > now_ms())
}
fn subscribe_authority_assignments(
&self,
) -> watch::Receiver<Option<NativeAuthorityAssignment>> {
self.shared.authority_assignment.subscribe()
}
#[cfg(feature = "managed-group-encryption")]
async fn publish_managed_room(&self, batch: Vec<NativeManagedRoomPublish>) -> Result<()> {
if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
bail!(
"managed room publish batch must contain 1..={MAX_MANAGED_BATCH_MESSAGES} messages"
);
}
self.request_unit(|reply| Command::PublishManaged { batch, reply })
.await
}
#[cfg(feature = "managed-group-encryption")]
fn subscribe_managed_room(&self) -> broadcast::Receiver<NativeManagedRoomMessage> {
self.shared.managed_messages.subscribe()
}
#[cfg(feature = "managed-group-encryption")]
fn avenue_key(&self) -> String {
format!("{}:{}", self.avenue.kind, self.avenue.id)
}
async fn request_unit(
&self,
build: impl FnOnce(oneshot::Sender<Result<()>>) -> Command,
) -> Result<()> {
let (reply, response) = oneshot::channel();
self.commands
.send(build(reply))
.await
.map_err(|_| anyhow!("coordination gateway actor stopped"))?;
response
.await
.map_err(|_| anyhow!("coordination gateway actor dropped its reply"))?
}
async fn mint_credential(
&self,
desired: &DesiredPresence,
purpose: &str,
refresh_grant: Option<String>,
) -> Result<NativeGatewayGrant> {
let ticket_fingerprint = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(Sha256::digest(desired.ticket.as_bytes()));
let credential = self
.grant_provider
.grant(NativeGatewayGrantRequest {
avenue: self.avenue.clone(),
device_id: self.device_id.clone(),
runtime_instance_id: self.runtime_instance_id.clone(),
ticket_fingerprint,
purpose: purpose.to_string(),
architecture: self.architecture,
room_delivery: self.room_delivery,
refresh_grant,
})
.await
.context("obtain native coordination grant")?;
if credential.protocol_version != GRANT_PROTOCOL_VERSION {
return Err(gateway_connect_error(
"unsupported native coordination protocol",
false,
));
}
let configured = reqwest::Url::parse(&self.endpoint)?;
let returned = reqwest::Url::parse(&credential.gateway_url)?;
if normalized_origin_scheme(configured.scheme())
!= normalized_origin_scheme(returned.scheme())
|| configured.host_str() != returned.host_str()
|| configured.port_or_known_default() != returned.port_or_known_default()
{
return Err(gateway_connect_error(
"native coordination credential returned an unexpected gateway origin",
false,
));
}
Ok(credential)
}
fn gateway_url(&self, credential: &NativeGatewayGrant) -> Result<String> {
let mut url = reqwest::Url::parse(&credential.gateway_url)?;
match url.scheme() {
"https" => url
.set_scheme("wss")
.map_err(|_| anyhow!("invalid gateway scheme"))?,
"http" => url
.set_scheme("ws")
.map_err(|_| anyhow!("invalid gateway scheme"))?,
"wss" | "ws" => {}
_ => bail!("native coordination gateway must use HTTPS/WSS"),
}
if url.query().is_some() || url.fragment().is_some() || !url.username().is_empty() {
bail!("native coordination gateway URL contains forbidden credentials or parameters");
}
let path = format!(
"{}/v{}/connect/{}",
url.path().trim_end_matches('/'),
credential.protocol_version,
credential.route_key
);
url.set_path(&path);
Ok(url.to_string())
}
fn device_from_gateway(&self, value: GatewayDevice) -> Device {
Device {
app_tag: Some(self.app_tag.clone()),
device_id: value.device_id,
user_id: value.user_id,
device_name: value.device_name,
platform_type: Some(value.platform_type),
capabilities: value.capabilities,
session_id: Some(value.runtime_instance_id),
node_id: Some(value.node_id),
tag: None,
kind: None,
metadata: value.metadata,
online: value.online,
ticket: Some(value.ticket),
last_seen_at: Some(serde_json::json!(value.updated_at_ms)),
expires_at: Some(serde_json::json!(value.expires_at_ms)),
created_at: None,
updated_at: Some(serde_json::json!(value.updated_at_ms)),
excluded_peers: value.excluded_peers,
}
}
fn member_from_gateway(&self, value: GatewayMember) -> Device {
Device {
app_tag: Some(self.app_tag.clone()),
device_id: value.device_id,
user_id: value.user_id,
device_name: value.device_name,
platform_type: Some(value.platform_type),
capabilities: value.capabilities,
session_id: None,
node_id: None,
tag: None,
kind: None,
metadata: value.metadata,
online: value.online,
ticket: None,
last_seen_at: Some(serde_json::json!(value.updated_at_ms)),
expires_at: Some(serde_json::json!(value.expires_at_ms)),
created_at: None,
updated_at: Some(serde_json::json!(value.updated_at_ms)),
excluded_peers: Vec::new(),
}
}
}
async fn rebuild_sparse_devices(adapter: &NativeCoordinationGatewaySignaling) {
let members = adapter.shared.members.read().await.clone();
let routes = adapter.shared.topology_routes.read().await.clone();
let mut next = HashMap::new();
for member in members.into_values() {
let device_id = member.device_id.clone();
let device = routes
.get(&device_id)
.cloned()
.map(|route| adapter.device_from_gateway(route))
.unwrap_or_else(|| adapter.member_from_gateway(member));
next.insert(device_id, device);
}
let mut current = adapter.shared.devices.write().await;
let mut events = current
.keys()
.filter(|device_id| !next.contains_key(*device_id))
.cloned()
.map(|device_id| DeviceEvent::Removed { device_id })
.collect::<Vec<_>>();
for (device_id, device) in &next {
match current.get(device_id) {
None => events.push(DeviceEvent::Added {
device: device.clone(),
}),
Some(existing) if existing != device => events.push(DeviceEvent::Modified {
device: device.clone(),
}),
Some(_) => {}
}
}
*current = next;
drop(current);
if !events.is_empty() {
let _ = adapter.shared.device_events.send(events);
}
}
type GatewaySocket =
tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
struct ConnectedGateway {
socket: GatewaySocket,
credential_expires_at_ms: u64,
credential_token: String,
lease_refresh_mode: LeaseRefreshMode,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LeaseRefreshMode {
Off,
Shadow,
Active,
}
impl LeaseRefreshMode {
fn from_wire(value: Option<&str>) -> Self {
match value {
Some("shadow") => Self::Shadow,
Some("active") => Self::Active,
_ => Self::Off,
}
}
}
async fn run_actor(
adapter: Arc<NativeCoordinationGatewaySignaling>,
mut commands: mpsc::Receiver<Command>,
) {
let mut socket: Option<ConnectedGateway> = None;
let mut reconnect_attempt = 0_u8;
let mut retry_at: Option<tokio::time::Instant> = None;
let mut pending_publish: Option<(DesiredPresence, oneshot::Sender<Result<()>>)> = None;
let mut circuit_failure: Option<(DesiredPresence, String)> = None;
loop {
if socket.is_none() {
if pending_publish
.as_ref()
.is_some_and(|(_, reply)| reply.is_closed())
{
pending_publish = None;
retry_at = None;
reconnect_attempt = 0;
*adapter.shared.desired.write().await = None;
*adapter.shared.applied_presence.write().await = None;
}
let retry_delay = retry_at
.map(|deadline| deadline.saturating_duration_since(tokio::time::Instant::now()))
.unwrap_or(Duration::from_secs(365 * 24 * 60 * 60));
tokio::select! {
command = commands.recv() => {
let Some(command) = command else { break; };
match command {
Command::Publish { desired, reply } => {
if let Some((blocked, message)) = circuit_failure.as_ref() {
if blocked == &desired {
let _ = reply.send(Err(anyhow!(message.clone())));
continue;
}
}
circuit_failure = None;
if let Some((_, previous_reply)) = pending_publish.take() {
let _ = previous_reply.send(Err(anyhow!(
"native coordination publication was superseded"
)));
}
*adapter.shared.desired.write().await = Some(desired.clone());
*adapter.shared.applied_presence.write().await = None;
pending_publish = Some((desired, reply));
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
}
Command::Stop => {
if let Some((_, reply)) = pending_publish.take() {
let _ = reply.send(Err(anyhow!(
"native coordination gateway stopped"
)));
}
break;
}
Command::Delete { user_id, device_id, reply } => {
let result = execute_control_delete(
&adapter,
&user_id,
&device_id,
)
.await;
let _ = reply.send(result);
}
other => reject_command(
other,
if retry_at.is_some() {
"native coordination gateway is reconnecting"
} else {
"native coordination gateway is not connected"
},
),
}
}
_ = tokio::time::sleep(retry_delay), if retry_at.is_some() => {
retry_at = None;
let desired = adapter.shared.desired.read().await.clone();
let Some(desired) = desired else {
reconnect_attempt = 0;
continue;
};
let result = connect(&adapter).await;
match result {
Ok(connected) => {
reconnect_attempt = 0;
circuit_failure = None;
*adapter.shared.applied_presence.write().await = Some(desired.clone());
socket = Some(connected);
if let Some((published, reply)) = pending_publish.take() {
if published == desired {
let _ = reply.send(Ok(()));
} else {
let _ = reply.send(Err(anyhow!(
"native coordination publication was superseded"
)));
}
}
}
Err(error) => {
reconnect_attempt = reconnect_attempt.saturating_add(1);
let retryable = is_retryable_gateway_error(&error);
eprintln!(
"[openrtc][coordination-gateway][connect-retry] attempt={}/{} retryable={} error={}",
reconnect_attempt,
MAX_RECONNECT_ATTEMPTS,
retryable,
error
);
if !retryable || reconnect_attempt >= MAX_RECONNECT_ATTEMPTS {
retry_at = None;
let failure = format!(
"native coordination gateway unavailable after {} attempts: {}",
reconnect_attempt, error
);
circuit_failure = Some((desired, failure.clone()));
if let Some((_, reply)) = pending_publish.take() {
let _ = reply.send(Err(anyhow!(failure)));
}
} else {
retry_at = Some(
tokio::time::Instant::now()
+ reconnect_delay(reconnect_attempt),
);
}
}
}
}
}
continue;
}
if let Some(active) = socket.as_mut() {
tokio::select! {
command = commands.recv() => {
let Some(command) = command else { break; };
if matches!(command, Command::Stop) {
let _ = active.socket.close(None).await;
break;
}
if let Some((blocked, message)) = circuit_failure.as_ref() {
if matches!(
&command,
Command::Publish { desired, .. } if desired == blocked
) {
let Command::Publish { reply, .. } = command else {
unreachable!("only unchanged publication enters this branch")
};
let _ = reply.send(Err(anyhow!(message.clone())));
continue;
}
}
let credential_rebind = match &command {
Command::Publish { desired, .. } => adapter
.shared
.applied_presence
.read()
.await
.as_ref()
.is_some_and(|applied| applied.ticket != desired.ticket),
_ => false,
};
if credential_rebind {
let Command::Publish { desired, reply } = command else {
unreachable!("credential rebind is only set for publication")
};
*adapter.shared.desired.write().await = Some(desired.clone());
*adapter.shared.applied_presence.write().await = None;
pending_publish = Some((desired, reply));
let _ = active.socket.close(None).await;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
continue;
}
let going_offline = matches!(&command, Command::Offline { .. });
let publishing = matches!(&command, Command::Publish { .. });
let deleting_local_device =
command_deletes_device(&command, &adapter.device_id);
let result = handle_command(&adapter, &mut active.socket, command).await;
if going_offline || (deleting_local_device && result.is_ok()) {
let _ = active.socket.close(None).await;
*adapter.shared.desired.write().await = None;
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = None;
circuit_failure = None;
} else if let Err(error) = result {
if is_retryable_gateway_error(&error) {
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
} else {
eprintln!(
"[openrtc][coordination-gateway][operation-rejected] retryable=false error={}",
error
);
if publishing {
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error.to_string()));
}
}
}
}
}
_ = tokio::time::sleep(auth_refresh_delay(active.credential_expires_at_ms)) => {
let refresh = if active.lease_refresh_mode == LeaseRefreshMode::Active {
refresh_lease(&adapter, &mut active.socket).await
} else {
refresh_authentication(
&adapter,
&mut active.socket,
Some(active.credential_token.clone()),
).await
};
match refresh {
Ok((expires_at_ms, credential_token)) => {
active.credential_expires_at_ms = expires_at_ms;
if let Some(credential_token) = credential_token {
active.credential_token = credential_token;
}
}
Err(error) => {
let terminal_error = if active.lease_refresh_mode == LeaseRefreshMode::Active {
eprintln!(
"[openrtc][coordination-gateway][lease-refresh-failed] retryable={} error={}",
is_retryable_gateway_error(&error),
error,
);
match refresh_authentication(
&adapter,
&mut active.socket,
None,
).await {
Ok((expires_at_ms, Some(token))) => {
active.credential_expires_at_ms = expires_at_ms;
active.credential_token = token;
continue;
}
Ok(_) => terminal_refresh_error(error, None),
Err(fallback) => {
eprintln!(
"[openrtc][coordination-gateway][lease-fallback-failed] retryable={} error={}",
is_retryable_gateway_error(&fallback),
fallback,
);
terminal_refresh_error(error, Some(fallback))
}
}
} else {
error
};
eprintln!(
"[openrtc][coordination-gateway][auth-refresh-failed] retryable={} error={}",
is_retryable_gateway_error(&terminal_error), terminal_error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
if is_retryable_gateway_error(&terminal_error) {
retry_at = Some(tokio::time::Instant::now());
} else {
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, terminal_error.to_string()));
}
}
}
}
}
_ = tokio::time::sleep(SOCKET_KEEPALIVE_INTERVAL) => {
if let Err(error) = send_socket_keepalive(&mut active.socket).await {
eprintln!(
"[openrtc][coordination-gateway][keepalive-failed] retryable=true error={}",
error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
}
}
incoming = active.socket.next() => {
match incoming {
Some(Ok(message)) => {
if let Err(error) = handle_message(&adapter, &mut active.socket, message).await {
eprintln!(
"[openrtc][coordination-gateway][socket-frame-failed] retryable={} error={:#}",
is_retryable_gateway_error(&error),
error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
if is_retryable_gateway_error(&error) {
retry_at = Some(tokio::time::Instant::now());
} else {
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error.to_string()));
}
}
}
}
Some(Err(error)) => {
eprintln!(
"[openrtc][coordination-gateway][socket-receive-failed] retryable=true error={:#}",
error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
}
None => {
eprintln!(
"[openrtc][coordination-gateway][socket-closed] retryable=true"
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now());
},
}
}
}
}
}
}
fn command_deletes_device(command: &Command, local_device_id: &str) -> bool {
matches!(
command,
Command::Delete { device_id, .. } if device_id == local_device_id
)
}
async fn connect(adapter: &NativeCoordinationGatewaySignaling) -> Result<ConnectedGateway> {
let desired = adapter
.shared
.desired
.read()
.await
.clone()
.ok_or_else(|| anyhow!("native coordination presence is not configured"))?;
connect_with_desired(adapter, &desired, "presence").await
}
async fn connect_with_desired(
adapter: &NativeCoordinationGatewaySignaling,
desired: &DesiredPresence,
purpose: &str,
) -> Result<ConnectedGateway> {
let credential = adapter.mint_credential(desired, purpose, None).await?;
if credential.expires_at_ms <= now_ms().saturating_add(30_000) {
bail!("native coordination credential expires too soon");
}
if purpose == "presence" {
*adapter.shared.topology_revision.write().await = 0;
*adapter.shared.active_grant.write().await = Some((
gateway_grant_jti(&credential.token)?,
credential.expires_at_ms,
));
}
let gateway_url = adapter.gateway_url(&credential)?;
let uri: Uri = gateway_url.parse().context("parse native gateway URI")?;
let mut request = uri.into_client_request()?;
request.headers_mut().insert(
"Sec-WebSocket-Protocol",
HeaderValue::from_str(&format!(
"{}, {}{}",
GATEWAY_PROTOCOL, GATEWAY_AUTH_PROTOCOL_PREFIX, credential.token
))?,
);
let (mut socket, _) =
tokio::time::timeout(REQUEST_TIMEOUT, tokio_tungstenite::connect_async(request))
.await
.context("native coordination WebSocket connect timed out")??;
socket
.send(Message::Text(
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "auth",
"token": credential.token,
"device": gateway_device_value(adapter, &desired).await
})
.to_string()
.into(),
))
.await?;
let lease_refresh_mode = tokio::time::timeout(ACK_TIMEOUT, async {
while let Some(message) = socket.next().await {
let message = message?;
if let Some(frame) = decode_frame(message)? {
match frame {
ServerFrame::Ready {
lease_refresh_mode, ..
} => {
return Ok::<LeaseRefreshMode, anyhow::Error>(LeaseRefreshMode::from_wire(
lease_refresh_mode.as_deref(),
))
}
ServerFrame::Error {
code,
message,
retryable,
..
} => {
return Err(gateway_connect_error(
format!("coordination gateway {code}: {message}"),
retryable,
));
}
other => handle_frame_with_socket(adapter, &mut socket, other).await?,
}
}
}
bail!("coordination gateway closed before ready")
})
.await
.context("native coordination authentication timed out")??;
Ok(ConnectedGateway {
socket,
credential_expires_at_ms: credential.expires_at_ms,
credential_token: credential.token,
lease_refresh_mode,
})
}
async fn execute_control_delete(
adapter: &NativeCoordinationGatewaySignaling,
user_id: &str,
target_device_id: &str,
) -> Result<()> {
let desired = DesiredPresence {
user_id: user_id.to_string(),
local_node_id: adapter.device_id.clone(),
ticket: format!("openrtc-device-control:{}", adapter.runtime_instance_id),
device_name: adapter.device_id.clone(),
metadata: None,
ttl_ms: 0,
online: false,
};
let mut connected = connect_with_desired(adapter, &desired, "device-control").await?;
let result = execute_operation(
adapter,
&mut connected.socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "device.delete",
"idempotencyKey": random_id("device-delete"),
"targetDeviceId": target_device_id,
}),
)
.await;
let _ = connected.socket.close(None).await;
result
}
async fn refresh_authentication(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
current_grant: Option<String>,
) -> Result<(u64, Option<String>)> {
let desired = adapter
.shared
.desired
.read()
.await
.clone()
.ok_or_else(|| anyhow!("native coordination presence is not configured"))?;
let credential = adapter
.mint_credential(&desired, "presence", current_grant)
.await?;
let refreshed_token = credential.token.clone();
let refreshed_jti = gateway_grant_jti(&credential.token)?;
socket
.send(Message::Text(
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "auth.refresh",
"token": credential.token,
})
.to_string()
.into(),
))
.await?;
tokio::time::timeout(ACK_TIMEOUT, async {
while let Some(message) = socket.next().await {
let message = message?;
if let Some(frame) = decode_frame(message)? {
match frame {
ServerFrame::AuthRefreshed { _expires_at_ms } => {
*adapter.shared.active_grant.write().await =
Some((refreshed_jti, _expires_at_ms));
update_local_device_expiry(adapter, _expires_at_ms).await;
return Ok((_expires_at_ms, Some(refreshed_token)));
}
ServerFrame::Error {
code,
message,
retryable,
..
} => {
return Err(gateway_connect_error(
format!("coordination gateway {code}: {message}"),
retryable,
));
}
other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
}
}
}
bail!("coordination gateway closed before authentication refresh")
})
.await
.context("native coordination authentication refresh timed out")?
}
async fn refresh_lease(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
) -> Result<(u64, Option<String>)> {
let idempotency_key = random_id("lease");
socket
.send(Message::Text(
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "lease.refresh",
"idempotencyKey": idempotency_key,
})
.to_string()
.into(),
))
.await?;
tokio::time::timeout(ACK_TIMEOUT, async {
while let Some(message) = socket.next().await {
let message = message?;
if let Some(frame) = decode_frame(message)? {
match frame {
ServerFrame::LeaseRefreshed {
idempotency_key: response_key,
expires_at_ms,
} if response_key == idempotency_key => {
update_local_device_expiry(adapter, expires_at_ms).await;
return Ok((expires_at_ms, None));
}
ServerFrame::Error {
code,
message,
retryable,
..
} => {
return Err(gateway_connect_error(
format!("coordination gateway {code}: {message}"),
retryable,
));
}
other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
}
}
}
bail!("coordination gateway closed before lease refresh acknowledgement")
})
.await
.context("native coordination lease refresh timed out")?
}
async fn update_local_device_expiry(
adapter: &NativeCoordinationGatewaySignaling,
expires_at_ms: u64,
) {
let event = {
let mut devices = adapter.shared.devices.write().await;
let Some(device) = devices.get_mut(&adapter.device_id) else {
return;
};
let next_expiry = serde_json::json!(expires_at_ms);
if device.expires_at.as_ref() == Some(&next_expiry) {
return;
}
device.expires_at = Some(next_expiry);
DeviceEvent::Modified {
device: device.clone(),
}
};
let _ = adapter.shared.device_events.send(vec![event]);
}
async fn handle_command(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
command: Command,
) -> Result<()> {
match command {
Command::Publish { desired, reply } => {
*adapter.shared.desired.write().await = Some(desired.clone());
if adapter.shared.applied_presence.read().await.as_ref() == Some(&desired) {
let _ = reply.send(Ok(()));
return Ok(());
}
let result =
execute_operation(adapter, socket, presence_frame(adapter, &desired).await).await;
if result.is_ok() {
*adapter.shared.applied_presence.write().await = Some(desired);
}
complete_command(reply, result, "native presence publication")?;
}
Command::Patch { patch, reply } => {
merge_patch(
&mut *adapter.shared.staged_patch.write().await,
patch.clone(),
);
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "device.patch",
"idempotencyKey": random_id("device-patch"),
"patch": patch,
}),
)
.await;
complete_command(reply, result, "native device patch")?;
}
Command::Offline { reply } => {
if let Some(desired) = adapter.shared.desired.write().await.as_mut() {
desired.online = false;
}
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "presence.offline",
"idempotencyKey": random_id("offline"),
}),
)
.await;
complete_command(reply, result, "native presence offline")?;
}
Command::Delete {
device_id, reply, ..
} => {
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "device.delete",
"idempotencyKey": random_id("device-delete"),
"targetDeviceId": device_id,
}),
)
.await;
complete_command(reply, result, "native device delete")?;
}
Command::SendSignal {
target_device_id,
payload,
state,
reply_payload,
reply,
} => {
let idempotency_key = random_id("signal");
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "signal.send",
"idempotencyKey": idempotency_key,
"targetDeviceId": target_device_id,
"payload": payload,
"state": state,
"replyPayload": reply_payload,
}),
)
.await
.map(|_| idempotency_key);
complete_command(reply, result, "native signal send")?;
}
Command::PutSession {
session_id,
session,
expires_at_ms,
reply,
} => {
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "session.put",
"idempotencyKey": random_id("session-put"),
"sessionId": session_id,
"session": session,
"expiresAtMs": expires_at_ms,
}),
)
.await;
complete_command(reply, result, "native session update")?;
}
Command::PublishAuthorityAssignment {
assignment,
signature,
reply,
} => {
let result = execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "authority.assignment.publish",
"idempotencyKey": random_id("authority-assignment"),
"assignment": assignment,
"signature": signature,
}),
)
.await;
complete_command(reply, result, "native authority assignment publish")?;
}
#[cfg(feature = "managed-group-encryption")]
Command::PublishManaged { batch, reply } => {
let result = publish_native_managed_batch(adapter, socket, batch).await;
complete_command(reply, result, "native managed room publish")?;
}
Command::Stop => {}
}
Ok(())
}
fn complete_command<T>(
reply: oneshot::Sender<Result<T>>,
result: Result<T>,
operation: &str,
) -> Result<()> {
match result {
Ok(value) => {
let _ = reply.send(Ok(value));
Ok(())
}
Err(error) => {
let retryable = is_retryable_gateway_error(&error);
let message = error.to_string();
let _ = reply.send(Err(anyhow!(message.clone())));
Err(gateway_connect_error(
format!("{operation} failed: {message}"),
retryable,
))
}
}
}
async fn execute_operation(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
frame: serde_json::Value,
) -> Result<()> {
let operation = frame
.get("type")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown")
.to_string();
let idempotency_key = frame
.get("idempotencyKey")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| anyhow!("gateway operation is missing idempotencyKey"))?
.to_string();
for attempt in 0..MAX_OPERATION_ATTEMPTS {
socket
.send(Message::Text(frame.to_string().into()))
.await
.context("send native gateway operation")?;
let renewal_required = tokio::time::timeout(ACK_TIMEOUT, async {
while let Some(message) = socket.next().await {
let message = message?;
if let Some(frame) = decode_frame(message)? {
match frame {
ServerFrame::Ack {
idempotency_key: ack,
} if ack == idempotency_key => return Ok(false),
ServerFrame::Error {
code,
message,
idempotency_key: Some(key),
retryable,
} if key == idempotency_key => {
if operation_requires_budget_renewal(&code, retryable) {
return Ok(true);
}
return Err(gateway_connect_error(
format!(
"coordination gateway {code}: {message} operation={operation}"
),
retryable,
));
}
other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
}
}
}
bail!("coordination gateway closed before acknowledgement")
})
.await
.context("native coordination acknowledgement timed out")??;
if !renewal_required {
return Ok(());
}
if attempt > 0 {
bail!("native coordination budget renewal did not fund operation={operation}");
}
refresh_authentication(adapter, socket, None)
.await
.with_context(|| format!("renew native coordination budget operation={operation}"))?;
}
unreachable!("native coordination operation loop returns or errors")
}
async fn handle_message(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
message: Message,
) -> Result<()> {
if let Some(frame) = decode_frame(message)? {
handle_frame_with_socket(adapter, socket, frame).await?;
}
Ok(())
}
#[cfg(feature = "managed-group-encryption")]
async fn ensure_managed_group(
adapter: &NativeCoordinationGatewaySignaling,
) -> Result<Option<ManagedGroupAction>> {
if adapter.shared.managed_group.lock().await.is_some() {
return Ok(None);
}
let signer = adapter
.managed_group_signer
.as_ref()
.ok_or_else(|| anyhow!("native managed rooms require a sign-only device key"))?;
let store = adapter
.managed_group_store
.as_ref()
.ok_or_else(|| anyhow!("native managed rooms require durable app-private state storage"))?;
let avenue_key = adapter.avenue_key();
let challenge = format!(
"openrtc:managed-group-wrapping-key:v1:{}:{}:{}",
adapter.app_tag, adapter.device_id, avenue_key,
);
let mut signature = signer
.sign(&adapter.app_tag, challenge.as_bytes())
.context("derive native managed group wrapping key")?;
if signature.len() < 32 {
bail!("native managed group device signature is invalid");
}
let mut digest = Sha256::new();
digest.update(b"openrtc:managed-group-key-derivation:v1:");
digest.update(challenge.as_bytes());
digest.update(&signature);
let wrapping_key: [u8; 32] = digest.finalize().into();
signature.zeroize();
let sealed_state = store
.load(&adapter.app_tag, &adapter.device_id, &avenue_key)
.await
.context("load native managed group state")?;
let controller = ManagedGroupController::new(
adapter.device_id.clone(),
wrapping_key,
sealed_state.as_deref(),
)?;
let action = controller.publish_key_package()?;
let mut current = adapter.shared.managed_group.lock().await;
if current.is_some() {
return Ok(None);
}
*current = Some(controller);
Ok(Some(action))
}
#[cfg(feature = "managed-group-encryption")]
async fn execute_managed_actions(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
actions: Vec<ManagedGroupAction>,
) -> Result<()> {
for action in actions {
match action {
ManagedGroupAction::PersistState { sealed_state } => {
let state = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(sealed_state)
.context("decode native managed group persisted state")?;
adapter
.managed_group_store
.as_ref()
.ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
.save(
&adapter.app_tag,
&adapter.device_id,
&adapter.avenue_key(),
&state,
)
.await
.context("persist native managed group state")?;
}
ManagedGroupAction::PublishKeyPackage {
key_package_hash,
key_package,
} => {
execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "managed.key-package.publish",
"idempotencyKey": random_id("managed-key-package"),
"keyPackageHash": key_package_hash,
"keyPackage": key_package,
}),
)
.await?;
}
ManagedGroupAction::RequestPreparePage {
preparation_id,
page_index,
} => {
execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "managed.prepare.page.request",
"idempotencyKey": random_id("managed-prepare-page"),
"preparationId": preparation_id,
"pageIndex": page_index,
}),
)
.await?;
}
ManagedGroupAction::SendArtifactChunk { chunk } => {
let mut frame = serde_json::to_value(chunk)?;
frame["v"] = serde_json::json!(WIRE_PROTOCOL_VERSION);
frame["type"] = serde_json::json!("managed.commit.chunk");
frame["idempotencyKey"] = serde_json::json!(random_id("managed-artifact"));
execute_operation(adapter, socket, frame).await?;
}
ManagedGroupAction::AcknowledgeCommit {
preparation_id,
encryption_epoch,
group_id_hash,
} => {
execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "managed.commit.ack",
"idempotencyKey": random_id("managed-commit-ack"),
"preparationId": preparation_id,
"encryptionEpoch": encryption_epoch,
"groupIdHash": group_id_hash,
}),
)
.await?;
}
}
}
Ok(())
}
#[cfg(feature = "managed-group-encryption")]
async fn ensure_managed_group_published(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
) -> Result<()> {
let Some(action) = ensure_managed_group(adapter).await? else {
return Ok(());
};
if let Err(error) = execute_managed_actions(adapter, socket, vec![action]).await {
*adapter.shared.managed_group.lock().await = None;
return Err(error).context("publish native managed group KeyPackage");
}
Ok(())
}
async fn handle_frame_with_socket(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
frame: ServerFrame,
) -> Result<()> {
#[cfg(feature = "managed-group-encryption")]
match frame {
ServerFrame::TopologyLease { lease } => {
accept_topology_lease(adapter, lease).await?;
if adapter
.room_architecture()
.is_some_and(|snapshot| snapshot.effective == EffectiveRoomArchitecture::Managed)
{
ensure_managed_group_published(adapter, socket).await?;
}
return Ok(());
}
ServerFrame::ManagedPreparePage { page } => {
ensure_managed_group_published(adapter, socket).await?;
let actions = adapter
.shared
.managed_group
.lock()
.await
.as_mut()
.ok_or_else(|| anyhow!("native managed group is unavailable"))?
.handle_prepare_page(page)?;
execute_managed_actions(adapter, socket, actions).await?;
return Ok(());
}
ServerFrame::ManagedArtifactChunk { chunk } => {
ensure_managed_group_published(adapter, socket).await?;
let actions = adapter
.shared
.managed_group
.lock()
.await
.as_mut()
.ok_or_else(|| anyhow!("native managed group is unavailable"))?
.handle_artifact_chunk(chunk)?;
execute_managed_actions(adapter, socket, actions).await?;
return Ok(());
}
ServerFrame::ManagedReceived {
sender_device_id,
architecture_epoch,
encryption_epoch,
batch,
} => {
handle_native_managed_delivery(
adapter,
sender_device_id,
architecture_epoch,
encryption_epoch,
batch,
)
.await?;
return Ok(());
}
other => return handle_frame(adapter, other).await,
}
#[cfg(not(feature = "managed-group-encryption"))]
{
let _ = socket;
handle_frame(adapter, frame).await
}
}
fn valid_reservation_id(value: &str) -> bool {
!value.is_empty()
&& value.len() <= 160
&& value.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b':' | b'@' | b'.')
})
}
async fn accept_topology_lease(
adapter: &NativeCoordinationGatewaySignaling,
lease: GatewayTopologyLease,
) -> Result<()> {
let invalid = || gateway_connect_error("native coordination topology lease is invalid", false);
let current_revision = *adapter.shared.topology_revision.read().await;
let active_grant = adapter.shared.active_grant.read().await.clone();
let Some((grant_jti, grant_expires_at_ms)) = active_grant else {
return Err(invalid());
};
if lease.schema_version != 2
|| lease.topology_revision == 0
|| lease.avenue != adapter.avenue
|| lease.grant_jti != grant_jti
|| lease.expires_at_ms <= now_ms()
|| lease.expires_at_ms > grant_expires_at_ms
{
return Err(invalid());
}
if lease.topology_revision <= current_revision {
return Ok(());
}
let current_architecture = adapter.shared.room_architecture.borrow().clone();
let current_encryption_epoch = *adapter.shared.encryption_epoch.read().await;
let routes = lease
.active
.iter()
.chain(&lease.backups)
.collect::<Vec<_>>();
let route_ids = routes
.iter()
.map(|route| route.device_id.as_str())
.collect::<std::collections::HashSet<_>>();
let members = adapter.shared.members.read().await;
let current_epoch = current_architecture
.as_ref()
.map_or(0, |snapshot| snapshot.epoch);
let architecture = &lease.architecture;
let valid_group_encryption = match lease.group_encryption.as_ref() {
None => {
architecture.effective_mode != EffectiveRoomArchitecture::Managed
|| architecture.phase == RoomArchitecturePhase::Preparing
}
Some(encryption) => {
encryption.policy_version == "room-mls-v1"
&& encryption.epoch > 0
&& encryption.epoch >= current_encryption_epoch
&& (encryption.epoch == current_encryption_epoch
|| current_encryption_epoch == 0
|| encryption.previous_encryption_epoch == Some(current_encryption_epoch))
&& matches!(encryption.phase.as_str(), "preparing" | "settled")
&& (architecture.effective_mode != EffectiveRoomArchitecture::Managed
|| architecture.phase != RoomArchitecturePhase::Settled
|| encryption.phase == "settled")
&& valid_reservation_id(&encryption.committer_device_id)
&& encryption.group_id_hash.len() == 43
&& encryption
.group_id_hash
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
}
};
let fixed_request_matches = adapter.architecture.is_none_or(|requested| {
requested == RoomArchitectureMode::Auto || requested == architecture.requested_mode
});
if (current_revision == 0 && lease.previous_topology_revision.is_some())
|| (current_revision > 0 && lease.previous_topology_revision != Some(current_revision))
|| lease.active.len() > 16
|| lease.backups.len() > 4
|| routes.len() > 20
|| route_ids.len() != routes.len()
|| routes.iter().any(|route| {
!route.online
|| route.device_id == adapter.device_id
|| route.node_id.trim().is_empty()
|| route.ticket.trim().is_empty()
|| !members.contains_key(&route.device_id)
})
|| architecture.policy_version != "room-architecture-v1"
|| architecture.epoch == 0
|| architecture.epoch < current_epoch
|| (architecture.epoch > current_epoch
&& current_epoch > 0
&& architecture.previous_architecture_epoch != Some(current_epoch))
|| !fixed_request_matches
|| architecture.held_credits_microusd > 9_007_199_254_740_991
|| architecture.quote_expires_at_ms == 0
|| architecture
.reservation_id
.as_deref()
.is_some_and(|value| !valid_reservation_id(value))
|| !valid_group_encryption
{
return Err(invalid());
}
drop(members);
let snapshot = RoomArchitectureSnapshot {
requested: architecture.requested_mode,
effective: architecture.effective_mode,
epoch: architecture.epoch,
phase: architecture.phase,
reason: architecture.reason,
held_credits_usd: architecture.held_credits_microusd as f64 / 1_000_000.0,
quote_expires_at_ms: architecture.quote_expires_at_ms,
};
let retains_authority_assignment = snapshot.effective == EffectiveRoomArchitecture::Authority;
adapter
.shared
.room_architecture
.send_if_modified(|current| {
if current.as_ref() == Some(&snapshot) {
false
} else {
*current = Some(snapshot);
true
}
});
if !retains_authority_assignment {
adapter.shared.authority_assignment.send_replace(None);
}
if let Some(encryption) = lease.group_encryption.as_ref() {
*adapter.shared.encryption_epoch.write().await = encryption.epoch;
}
*adapter.shared.topology_routes.write().await = lease
.active
.into_iter()
.chain(lease.backups)
.map(|route| (route.device_id.clone(), route))
.collect();
*adapter.shared.topology_revision.write().await = lease.topology_revision;
rebuild_sparse_devices(adapter).await;
Ok(())
}
fn decode_frame(message: Message) -> Result<Option<ServerFrame>> {
match message {
Message::Text(text) => Ok(Some(serde_json::from_str(&text)?)),
Message::Binary(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
Message::Ping(_) | Message::Pong(_) => Ok(None),
Message::Close(_) => bail!("coordination gateway closed"),
_ => Ok(None),
}
}
fn gateway_grant_jti(token: &str) -> Result<String> {
let payload = token
.split('.')
.nth(1)
.ok_or_else(|| anyhow!("native gateway grant is malformed"))?;
let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(payload)
.context("decode native gateway grant")?;
let value: serde_json::Value = serde_json::from_slice(&decoded)?;
let jti = value
.get("jti")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
if jti.is_empty()
|| jti.len() > 200
|| !jti.bytes().all(|byte| {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
})
{
bail!("native gateway grant jti is invalid");
}
Ok(jti.to_string())
}
async fn handle_frame(
adapter: &NativeCoordinationGatewaySignaling,
frame: ServerFrame,
) -> Result<()> {
match frame {
ServerFrame::RosterSnapshot { devices } => {
let mut mapped = HashMap::new();
let mut events = Vec::with_capacity(devices.len());
for raw in devices {
let device = adapter.device_from_gateway(raw);
mapped.insert(device.device_id.clone(), device.clone());
events.push(DeviceEvent::Added { device });
}
*adapter.shared.devices.write().await = mapped;
if !events.is_empty() {
let _ = adapter.shared.device_events.send(events);
}
}
ServerFrame::MembershipSnapshot { members } => {
*adapter
.shared
.membership_pages
.lock()
.map_err(|_| anyhow!("native membership page lock poisoned"))? = None;
*adapter.shared.members.write().await = members
.into_iter()
.map(|member| (member.device_id.clone(), member))
.collect();
rebuild_sparse_devices(adapter).await;
}
ServerFrame::MembershipPage {
snapshot_id,
page_index,
page_count,
member_count,
members,
} => {
let valid_snapshot_id = !snapshot_id.is_empty()
&& snapshot_id.len() <= 160
&& snapshot_id
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b':' | b'_' | b'-'));
let valid_envelope = valid_snapshot_id
&& page_count >= 2
&& page_count <= 50
&& page_index < page_count
&& member_count > 100
&& member_count <= 5_000
&& !members.is_empty()
&& members.len() <= 100;
let completed = {
let mut pending = adapter
.shared
.membership_pages
.lock()
.map_err(|_| anyhow!("native membership page lock poisoned"))?;
if page_index == 0 && valid_envelope {
*pending = Some(PendingMembershipPages {
snapshot_id: snapshot_id.clone(),
next_page_index: 0,
page_count,
member_count,
members: HashMap::new(),
});
}
let page = pending
.as_mut()
.ok_or_else(|| anyhow!("native membership page arrived without a snapshot"))?;
if !valid_envelope
|| page.snapshot_id != snapshot_id
|| page.next_page_index != page_index
|| page.page_count != page_count
|| page.member_count != member_count
|| members.iter().any(|member| {
member.device_id.is_empty() || page.members.contains_key(&member.device_id)
})
{
bail!("native membership page is invalid");
}
for member in members {
page.members.insert(member.device_id.clone(), member);
}
page.next_page_index += 1;
if page.next_page_index == page.page_count {
if page.members.len() != page.member_count {
bail!("native membership pages are incomplete");
}
pending.take().map(|complete| complete.members)
} else {
None
}
};
if let Some(members) = completed {
*adapter.shared.members.write().await = members;
rebuild_sparse_devices(adapter).await;
}
}
ServerFrame::MembershipChanged { operation, member } => {
let mut members = adapter.shared.members.write().await;
if operation == "delete" {
members.remove(&member.device_id);
} else {
members.insert(member.device_id.clone(), member);
}
drop(members);
rebuild_sparse_devices(adapter).await;
}
ServerFrame::TopologyLease { lease } => {
accept_topology_lease(adapter, lease).await?;
}
ServerFrame::PresenceChanged { operation, device } => {
let device = adapter.device_from_gateway(device);
let event = if operation == "delete" {
adapter
.shared
.devices
.write()
.await
.remove(&device.device_id);
DeviceEvent::Removed {
device_id: device.device_id,
}
} else {
let existed = adapter
.shared
.devices
.write()
.await
.insert(device.device_id.clone(), device.clone())
.is_some();
if existed {
DeviceEvent::Modified { device }
} else {
DeviceEvent::Added { device }
}
};
let _ = adapter.shared.device_events.send(vec![event]);
}
ServerFrame::SessionChanged {
operation,
session_id,
session,
} => {
let event = if operation == "delete" {
SessionEvent::Removed { session_id }
} else {
let mut value = session.unwrap_or_else(|| serde_json::json!({}));
value["connectionId"] = serde_json::Value::String(session_id);
SessionEvent::Modified {
session: serde_json::from_value(value)?,
}
};
let _ = adapter.shared.session_events.send(vec![event]);
}
ServerFrame::SignalReceived {
sender_device_id,
payload,
state,
reply_payload,
..
} => {
let mut pending = adapter
.shared
.pending_messages
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if pending.len() >= 1_000 {
pending.remove(0);
}
pending.push(SignalingEnvelope {
app_tag: Some(adapter.app_tag.clone()),
sender_id: sender_device_id,
target_id: adapter.device_id.clone(),
payload,
state,
reply_payload,
timestamp: now_ms() as i64,
sender_user_id: None,
target_user_id: None,
expires_at: None,
});
}
ServerFrame::Error {
code,
message,
retryable,
..
} => {
return Err(gateway_connect_error(
format!("coordination gateway {code}: {message}"),
retryable,
))
}
ServerFrame::AuthorityAssignment {
assignment,
signature,
} => {
validate_authority_assignment(&assignment, &signature)?;
if assignment.subject_device_id != adapter.device_id
|| !matches!(
adapter.shared.room_architecture.borrow().as_ref(),
Some(snapshot) if snapshot.effective == EffectiveRoomArchitecture::Authority
)
{
bail!("native authority assignment does not match the active room lease");
}
let next = NativeAuthorityAssignment {
assignment,
signature,
};
if let Some(current) = adapter.shared.authority_assignment.borrow().clone() {
let current_key = (current.assignment.generation, current.assignment.revision);
let next_key = (next.assignment.generation, next.assignment.revision);
if next_key < current_key || (next_key == current_key && next != current) {
bail!("native authority assignment generation is stale");
}
if next == current {
return Ok(());
}
}
adapter.shared.authority_assignment.send_replace(Some(next));
}
#[cfg(feature = "managed-group-encryption")]
ServerFrame::ManagedPreparePage { .. }
| ServerFrame::ManagedArtifactChunk { .. }
| ServerFrame::ManagedReceived { .. } => {
bail!("native managed room frame requires the active gateway owner")
}
ServerFrame::Ready { .. }
| ServerFrame::AuthRefreshed { .. }
| ServerFrame::LeaseRefreshed { .. }
| ServerFrame::Ack { .. }
| ServerFrame::Pong => {}
}
Ok(())
}
async fn presence_frame(
adapter: &NativeCoordinationGatewaySignaling,
desired: &DesiredPresence,
) -> serde_json::Value {
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "presence.upsert",
"idempotencyKey": random_id("presence"),
"ttlMs": desired.ttl_ms,
"device": gateway_device_value(adapter, desired).await,
})
}
async fn gateway_device_value(
adapter: &NativeCoordinationGatewaySignaling,
desired: &DesiredPresence,
) -> serde_json::Value {
let patch = adapter.shared.staged_patch.read().await.clone();
serde_json::to_value(OutboundGatewayDevice {
device_id: adapter.device_id.clone(),
runtime_instance_id: adapter.runtime_instance_id.clone(),
node_id: desired.local_node_id.clone(),
device_name: patch
.device_name
.unwrap_or_else(|| desired.device_name.clone()),
platform_type: adapter.platform_type.clone(),
ticket: desired.ticket.clone(),
metadata: patch.metadata.or_else(|| desired.metadata.clone()),
capabilities: patch.capabilities,
excluded_peers: patch.excluded_peers.unwrap_or_default(),
online: desired.online,
})
.expect("gateway device projection is serializable")
}
fn merge_patch(target: &mut DevicePatch, patch: DevicePatch) {
if patch.device_name.is_some() {
target.device_name = patch.device_name;
}
if patch.capabilities.is_some() {
target.capabilities = patch.capabilities;
}
if patch.metadata.is_some() {
target.metadata = patch.metadata;
}
if patch.excluded_peers.is_some() {
target.excluded_peers = patch.excluded_peers;
}
}
fn reject_command(command: Command, message: &str) {
match command {
Command::Publish { reply, .. }
| Command::Patch { reply, .. }
| Command::Offline { reply }
| Command::Delete { reply, .. }
| Command::PutSession { reply, .. } => {
let _ = reply.send(Err(anyhow!(message.to_string())));
}
Command::PublishAuthorityAssignment { reply, .. } => {
let _ = reply.send(Err(anyhow!(message.to_string())));
}
#[cfg(feature = "managed-group-encryption")]
Command::PublishManaged { reply, .. } => {
let _ = reply.send(Err(anyhow!(message.to_string())));
}
Command::SendSignal { reply, .. } => {
let _ = reply.send(Err(anyhow!(message.to_string())));
}
Command::Stop => {}
}
}
#[cfg(feature = "managed-group-encryption")]
async fn publish_native_managed_batch(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
batch: Vec<NativeManagedRoomPublish>,
) -> Result<()> {
if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
bail!("managed room publish batch exceeds its protected bound");
}
let architecture = adapter
.room_architecture()
.ok_or_else(|| anyhow!("native managed room architecture is unavailable"))?;
let encryption_epoch = *adapter.shared.encryption_epoch.read().await;
if architecture.effective != EffectiveRoomArchitecture::Managed
|| architecture.phase != RoomArchitecturePhase::Settled
|| encryption_epoch == 0
{
bail!("native managed room encryption is not settled");
}
let (envelopes, sealed_state) = {
let mut managed = adapter.shared.managed_group.lock().await;
let controller = managed
.as_mut()
.ok_or_else(|| anyhow!("native managed room group is unavailable"))?;
let mut envelopes = Vec::with_capacity(batch.len());
let mut sealed_state = None;
for message in batch {
let protected = controller.seal_payload(
architecture.epoch,
encryption_epoch,
&message.message_id,
&message.channel,
message.priority,
message.zone_id.as_deref(),
&message.payload,
)?;
sealed_state = Some(protected.sealed_state);
envelopes.push(NativeManagedEncryptedEnvelope {
message_id: message.message_id,
channel: message.channel,
priority: message.priority,
scope: message
.zone_id
.map_or(NativeManagedScope::Global, |zone_id| {
NativeManagedScope::Zone { zone_id }
}),
ciphertext: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(protected.data),
});
}
(
envelopes,
sealed_state.expect("non-empty managed batch has sealed state"),
)
};
adapter
.managed_group_store
.as_ref()
.ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
.save(
&adapter.app_tag,
&adapter.device_id,
&adapter.avenue_key(),
&sealed_state,
)
.await
.context("persist native managed group state before publish")?;
execute_operation(
adapter,
socket,
serde_json::json!({
"v": WIRE_PROTOCOL_VERSION,
"type": "managed.publish",
"idempotencyKey": random_id("managed-publish"),
"architectureEpoch": architecture.epoch,
"encryptionEpoch": encryption_epoch,
"batch": envelopes,
}),
)
.await
}
#[cfg(feature = "managed-group-encryption")]
async fn handle_native_managed_delivery(
adapter: &NativeCoordinationGatewaySignaling,
sender_device_id: String,
architecture_epoch: u64,
encryption_epoch: u64,
batch: Vec<NativeManagedEncryptedEnvelope>,
) -> Result<()> {
if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
bail!("native managed delivery batch exceeds its protected bound");
}
let architecture = adapter
.room_architecture()
.ok_or_else(|| anyhow!("native managed room architecture is unavailable"))?;
if architecture.effective != EffectiveRoomArchitecture::Managed
|| architecture.phase != RoomArchitecturePhase::Settled
|| architecture.epoch != architecture_epoch
|| *adapter.shared.encryption_epoch.read().await != encryption_epoch
{
bail!("native managed room delivery uses a stale lease");
}
let (messages, sealed_state) = {
let mut managed = adapter.shared.managed_group.lock().await;
let controller = managed
.as_mut()
.ok_or_else(|| anyhow!("native managed room group is unavailable"))?;
let mut messages = Vec::with_capacity(batch.len());
let mut sealed_state = None;
for envelope in batch {
let zone_id = match envelope.scope {
NativeManagedScope::Global => None,
NativeManagedScope::Zone { zone_id } => Some(zone_id),
};
let ciphertext = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(envelope.ciphertext)
.context("decode native managed ciphertext")?;
let protected = controller.open_payload(
architecture_epoch,
encryption_epoch,
&envelope.message_id,
&envelope.channel,
envelope.priority,
zone_id.as_deref(),
&ciphertext,
)?;
sealed_state = Some(protected.sealed_state);
messages.push(NativeManagedRoomMessage {
sender_device_id: sender_device_id.clone(),
message_id: envelope.message_id,
channel: envelope.channel,
priority: envelope.priority,
zone_id,
payload: protected.data,
});
}
(
messages,
sealed_state.expect("non-empty managed delivery has sealed state"),
)
};
adapter
.managed_group_store
.as_ref()
.ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
.save(
&adapter.app_tag,
&adapter.device_id,
&adapter.avenue_key(),
&sealed_state,
)
.await
.context("persist native managed group state before delivery")?;
for message in messages {
let _ = adapter.shared.managed_messages.send(message);
}
Ok(())
}
fn reconnect_delay(attempt: u8) -> Duration {
let exponent = u32::from(attempt.saturating_sub(1).min(7));
Duration::from_millis(
(250_u64.saturating_mul(1_u64 << exponent)).min(MAX_RECONNECT_DELAY.as_millis() as u64),
)
}
fn auth_refresh_delay(expires_at_ms: u64) -> Duration {
Duration::from_millis(
expires_at_ms
.saturating_sub(now_ms())
.saturating_sub(AUTH_REFRESH_SKEW_MS)
.max(1_000),
)
}
async fn send_socket_keepalive<S, E>(socket: &mut S) -> Result<()>
where
S: Sink<Message, Error = E> + Unpin,
E: StdError + Send + Sync + 'static,
{
socket
.send(Message::Ping(Vec::new().into()))
.await
.context("send native coordination WebSocket keepalive")
}
fn required(name: &str, value: String) -> Result<String> {
let value = value.trim().to_string();
if value.is_empty() {
bail!("native coordination {name} is required");
}
Ok(value)
}
fn validate_endpoint(name: &str, value: String, allow_websocket: bool) -> Result<String> {
let value = required(name, value)?;
let url = reqwest::Url::parse(&value)?;
let allowed = matches!(url.scheme(), "https" | "http")
|| (allow_websocket && matches!(url.scheme(), "wss" | "ws"));
let secure = matches!(url.scheme(), "https" | "wss");
let loopback = matches!(
url.host_str(),
Some("127.0.0.1") | Some("localhost") | Some("::1")
);
if !allowed
|| (!secure && !loopback)
|| !url.username().is_empty()
|| url.password().is_some()
|| url.query().is_some()
|| url.fragment().is_some()
{
bail!("native coordination {name} is invalid");
}
Ok(value.trim_end_matches('/').to_string())
}
fn normalized_origin_scheme(scheme: &str) -> &str {
match scheme {
"ws" => "http",
"wss" => "https",
other => other,
}
}
fn random_id(prefix: &str) -> String {
format!("{prefix}:{}", Uuid::new_v4())
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64
}
#[async_trait]
impl SignalingBackend for NativeCoordinationGatewaySignaling {
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<()> {
if !is_online {
merge_patch(
&mut *self.shared.staged_patch.write().await,
DevicePatch {
device_name: Some(name.to_string()),
metadata: metadata.map(str::to_string),
..DevicePatch::default()
},
);
return Ok(());
}
let local_node_id = iroh::EndpointId::from_str(local_node_id)
.context("native coordination node ID is invalid")?
.to_string();
let desired = DesiredPresence {
user_id: user_id.to_string(),
local_node_id,
ticket: ticket_str.to_string(),
device_name: name.to_string(),
metadata: metadata.map(str::to_string),
ttl_ms,
online: is_online,
};
*self.shared.desired.write().await = Some(desired.clone());
self.request_unit(|reply| Command::Publish { desired, reply })
.await
}
async fn set_offline(&self, _user_id: &str, _local_node_id: &str) -> Result<()> {
self.request_unit(|reply| Command::Offline { reply }).await
}
async fn update_live_presence(
&self,
user_id: &str,
local_node_id: &str,
ticket_str: &str,
name: &str,
metadata: Option<&str>,
) -> Result<()> {
self.update_presence(
user_id,
local_node_id,
ticket_str,
true,
name,
15 * 60_000,
metadata,
)
.await
}
async fn set_live_presence_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
self.set_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<()> {
if device_id != self.device_id {
bail!("native coordination socket can update only its local device");
}
let patch = DevicePatch {
device_name: device_name.map(str::to_string),
capabilities,
metadata: metadata.map(str::to_string),
excluded_peers: None,
};
merge_patch(&mut *self.shared.staged_patch.write().await, patch.clone());
if self.shared.applied_presence.read().await.is_none() {
return Ok(());
}
self.request_unit(|reply| Command::Patch { patch, reply })
.await
}
async fn delete_device(&self, _user_id: &str, device_id: &str) -> Result<()> {
self.request_unit(|reply| Command::Delete {
user_id: _user_id.to_string(),
device_id: device_id.to_string(),
reply,
})
.await
}
async fn set_excluded_peers(
&self,
_user_id: &str,
_local_node_id: &str,
excluded_peers: &[String],
) -> Result<()> {
let mut normalized = excluded_peers
.iter()
.map(|value| value.trim().to_ascii_lowercase())
.filter(|value| !value.is_empty())
.collect::<Vec<_>>();
normalized.sort();
normalized.dedup();
if normalized.len() > MAX_EXCLUDED_PEERS {
bail!(
"native coordination excluded-peer set exceeds {} devices",
MAX_EXCLUDED_PEERS
);
}
let patch = DevicePatch {
excluded_peers: Some(normalized),
..DevicePatch::default()
};
merge_patch(&mut *self.shared.staged_patch.write().await, patch.clone());
if self.shared.applied_presence.read().await.is_none() {
return Ok(());
}
self.request_unit(|reply| Command::Patch { patch, reply })
.await
}
async fn search_devices(
&self,
_user_id: &str,
exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
self.list_devices("", exclude_node_id).await
}
async fn list_devices(
&self,
_user_id: &str,
exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
let mut devices = self
.shared
.devices
.read()
.await
.values()
.filter(|device| exclude_node_id != device.node_id.as_deref())
.cloned()
.collect::<Vec<_>>();
devices.sort_by(|left, right| left.device_id.cmp(&right.device_id));
Ok(devices)
}
async fn send_message(
&self,
_sender_id: &str,
target_id: &str,
payload: &str,
state: Option<&str>,
reply_payload: Option<&str>,
) -> Result<String> {
let (reply, response) = oneshot::channel();
self.commands
.send(Command::SendSignal {
target_device_id: target_id.to_string(),
payload: payload.to_string(),
state: state.map(str::to_string),
reply_payload: reply_payload.map(str::to_string),
reply,
})
.await
.map_err(|_| anyhow!("coordination gateway actor stopped"))?;
response
.await
.map_err(|_| anyhow!("coordination gateway actor dropped its reply"))?
}
async fn subscribe_devices(
&self,
_user_id: &str,
) -> Result<BoxStream<'static, Result<Vec<DeviceEvent>>>> {
let receiver = self.shared.device_events.subscribe();
Ok(Box::pin(
tokio_stream::wrappers::BroadcastStream::new(receiver)
.filter_map(|event| async move { event.ok().map(Ok) }),
))
}
async fn create_session(&self, session: SignalingSession) -> Result<()> {
let expires_at_ms = session
.expires_at
.unwrap_or_else(|| now_ms() as i64 + 15 * 60_000);
let session_id = session.connection_id.clone();
let value = serde_json::to_value(session)?;
self.request_unit(|reply| Command::PutSession {
session_id,
session: value,
expires_at_ms,
reply,
})
.await
}
async fn update_session(&self, session_id: &str, update_data: serde_json::Value) -> Result<()> {
let expires_at_ms = update_data
.get("expiresAt")
.and_then(serde_json::Value::as_i64)
.unwrap_or_else(|| now_ms() as i64 + 15 * 60_000);
self.request_unit(|reply| Command::PutSession {
session_id: session_id.to_string(),
session: update_data,
expires_at_ms,
reply,
})
.await
}
async fn subscribe_sessions(
&self,
_local_device_id: &str,
) -> Result<BoxStream<'static, Result<Vec<SessionEvent>>>> {
let receiver = self.shared.session_events.subscribe();
Ok(Box::pin(
tokio_stream::wrappers::BroadcastStream::new(receiver)
.filter_map(|event| async move { event.ok().map(Ok) }),
))
}
}
impl Drop for NativeCoordinationGatewaySignaling {
fn drop(&mut self) {
let _ = self.commands.try_send(Command::Stop);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[derive(Default)]
struct RecordingGrantProvider {
requests: Mutex<Vec<NativeGatewayGrantRequest>>,
}
#[async_trait]
impl NativeGatewayGrantProvider for RecordingGrantProvider {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
self.requests.lock().expect("requests").push(request);
Ok(NativeGatewayGrant {
protocol_version: 2,
gateway_url: "https://gateway.example.test".to_string(),
route_key: "route-1".to_string(),
token: "grant-2".to_string(),
expires_at_ms: now_ms() + 3_600_000,
})
}
}
#[test]
fn native_authority_targets_match_the_public_assignment_bounds() {
let priority = (0..8)
.map(|index| format!("peer-{index}"))
.collect::<Vec<_>>();
let relevant = (0..32)
.map(|index| format!("entity-{index}"))
.collect::<Vec<_>>();
validate_native_authority_assignment_targets(&priority, &relevant)
.expect("the exact public bounds are accepted");
let mut too_many_priority = priority.clone();
too_many_priority.push("peer-8".to_string());
assert!(
validate_native_authority_assignment_targets(&too_many_priority, &relevant).is_err()
);
let mut too_many_relevant = relevant.clone();
too_many_relevant.push("entity-32".to_string());
assert!(
validate_native_authority_assignment_targets(&priority, &too_many_relevant).is_err()
);
assert!(validate_native_authority_assignment_targets(
&["peer-1".to_string(), "peer-1".to_string()],
&[],
)
.is_err());
assert!(validate_native_authority_assignment_targets(
&[],
&["entity-1".to_string(), "entity-1".to_string()],
)
.is_err());
}
#[tokio::test]
async fn native_capability_handle_is_inactive_until_runtime_presence_and_closes_once() {
let provider = Arc::new(RecordingGrantProvider::default());
let capabilities = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
"device-native",
"native",
provider.clone(),
)
.expect("capability namespace")
.with_testing_endpoint("https://gateway.example.test")
.expect("testing endpoint");
let handle = capabilities
.join_room_with_architecture("room-1", RoomArchitectureMode::Auto)
.expect("capability handle");
assert_eq!(handle.kind(), NativeCapabilityKind::Room);
assert_eq!(handle.id(), "room-1");
assert!(!handle.is_closed());
assert!(provider.requests.lock().expect("requests").is_empty());
handle.close().await;
handle.close().await;
assert!(handle.is_closed());
assert!(provider.requests.lock().expect("requests").is_empty());
}
#[tokio::test]
async fn room_architecture_intent_reaches_the_native_grant_owner() {
let provider = Arc::new(RecordingGrantProvider::default());
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-1".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "room".to_string(),
id: "room-1".to_string(),
},
architecture: Some(RoomArchitectureMode::Managed),
room_delivery: Some(RoomDelivery::Reliable),
grant_provider: provider.clone(),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
let desired = DesiredPresence {
user_id: "room-1".to_string(),
local_node_id: "node-1".to_string(),
ticket: "ticket-1".to_string(),
device_name: "Native".to_string(),
metadata: None,
ttl_ms: 900_000,
online: true,
};
adapter
.mint_credential(&desired, "presence", None)
.await
.expect("grant");
let requests = provider.requests.lock().expect("requests");
assert_eq!(requests.len(), 1);
assert_eq!(
requests[0].architecture,
Some(RoomArchitectureMode::Managed)
);
assert_eq!(requests[0].room_delivery, Some(RoomDelivery::Reliable));
}
#[tokio::test]
async fn native_authority_assignment_is_local_monotonic_and_bounded() {
let provider = Arc::new(RecordingGrantProvider::default());
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-local".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "room".to_string(),
id: "room-authority".to_string(),
},
architecture: Some(RoomArchitectureMode::Authority),
room_delivery: Some(RoomDelivery::Reliable),
grant_provider: provider,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
adapter
.shared
.room_architecture
.send_replace(Some(RoomArchitectureSnapshot {
requested: RoomArchitectureMode::Authority,
effective: EffectiveRoomArchitecture::Authority,
epoch: 2,
phase: RoomArchitecturePhase::Settled,
reason: RoomArchitectureReason::Manual,
held_credits_usd: 0.01,
quote_expires_at_ms: now_ms() + 30_000,
}));
let assignment = NativeAuthorityInterestAssignment {
policy_version: "room-authority-assignment-v1".to_string(),
service_id: "service-primary".to_string(),
generation: 2,
revision: 1,
subject_device_id: "device-local".to_string(),
shard_id: "zone-a".to_string(),
priority_device_ids: vec!["device-nearby".to_string()],
relevant_entity_ids: vec!["entity-player-1".to_string()],
issued_at_ms: now_ms(),
expires_at_ms: now_ms() + 60_000,
};
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([0_u8; 64]);
for _ in 0..2 {
handle_frame(
&adapter,
ServerFrame::AuthorityAssignment {
assignment: assignment.clone(),
signature: signature.clone(),
},
)
.await
.expect("current assignment");
}
assert_eq!(
adapter.authority_assignment().unwrap().assignment,
assignment
);
let stale = NativeAuthorityInterestAssignment {
generation: 1,
..assignment.clone()
};
assert!(handle_frame(
&adapter,
ServerFrame::AuthorityAssignment {
assignment: stale,
signature: signature.clone(),
},
)
.await
.is_err());
let wrong_subject = NativeAuthorityInterestAssignment {
revision: 2,
subject_device_id: "device-other".to_string(),
..assignment
};
assert!(handle_frame(
&adapter,
ServerFrame::AuthorityAssignment {
assignment: wrong_subject,
signature,
},
)
.await
.is_err());
}
#[test]
fn reconnect_is_bounded_and_capped() {
assert_eq!(reconnect_delay(1), Duration::from_millis(250));
assert_eq!(reconnect_delay(12), MAX_RECONNECT_DELAY);
assert_eq!(MAX_RECONNECT_ATTEMPTS, 12);
}
#[test]
fn endpoints_reject_query_credentials() {
assert!(validate_endpoint(
"endpoint",
"https://gateway.example.test?token=secret".to_string(),
true,
)
.is_err());
assert!(validate_endpoint(
"endpoint",
"https://user:secret@gateway.example.test".to_string(),
true,
)
.is_err());
assert!(
validate_endpoint("endpoint", "http://gateway.example.test".to_string(), true,)
.is_err()
);
assert!(validate_endpoint("endpoint", "http://127.0.0.1:8787".to_string(), true,).is_ok());
assert_eq!(normalized_origin_scheme("wss"), "https");
assert_eq!(normalized_origin_scheme("ws"), "http");
}
#[tokio::test]
async fn grant_provider_owns_native_issue_and_in_place_refresh() {
let provider = Arc::new(RecordingGrantProvider::default());
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-1".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "user".to_string(),
id: "principal-1".to_string(),
},
architecture: None,
room_delivery: None,
grant_provider: provider.clone(),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
let desired = DesiredPresence {
user_id: "principal-1".to_string(),
local_node_id: "node-1".to_string(),
ticket: "ticket-1".to_string(),
device_name: "Native".to_string(),
metadata: None,
ttl_ms: 900_000,
online: true,
};
let issued = adapter
.mint_credential(&desired, "presence", None)
.await
.expect("issue");
let refreshed = adapter
.mint_credential(&desired, "presence", Some(issued.token.clone()))
.await
.expect("refresh");
assert_eq!(issued.protocol_version, 2);
assert_eq!(refreshed.protocol_version, 2);
assert!(adapter
.gateway_url(&issued)
.expect("url")
.contains("/v2/connect/route-1"));
let requests = provider.requests.lock().expect("requests");
assert_eq!(requests.len(), 2);
assert_eq!(requests[0].refresh_grant, None);
assert_eq!(requests[1].refresh_grant.as_deref(), Some("grant-2"));
assert_eq!(requests[1].avenue.kind, "user");
assert_eq!(requests[1].avenue.id, "principal-1");
}
#[test]
fn patch_merge_preserves_unrelated_fields() {
let mut target = DevicePatch {
device_name: Some("old".into()),
capabilities: Some(DeviceCapabilities {
can_host: true,
can_sync: true,
read_only: false,
}),
metadata: None,
excluded_peers: None,
};
merge_patch(
&mut target,
DevicePatch {
excluded_peers: Some(vec!["peer".into()]),
..DevicePatch::default()
},
);
assert_eq!(target.device_name.as_deref(), Some("old"));
assert_eq!(target.excluded_peers, Some(vec!["peer".into()]));
}
#[test]
fn permanent_gateway_errors_are_not_retried() {
let permanent = gateway_connect_error("invalid payload", false);
let transient = gateway_connect_error("provider unavailable", true);
assert!(!is_retryable_gateway_error(&permanent));
assert!(is_retryable_gateway_error(&transient));
}
#[test]
fn only_typed_retryable_budget_exhaustion_renews_an_operation() {
assert!(operation_requires_budget_renewal(
"budget-renewal-required",
true,
));
assert!(!operation_requires_budget_renewal(
"budget-renewal-required",
false,
));
assert!(!operation_requires_budget_renewal(
"provider-budget-exhausted",
true,
));
}
#[test]
fn terminal_control_plane_rejections_are_not_retried() {
let revoked = anyhow::Error::new(ControlPlaneHttpError {
status: 403,
message: "device revoked".into(),
reason: Some("device-certificate-revoked".into()),
})
.context("obtain native coordination grant");
let throttled = anyhow::Error::new(ControlPlaneHttpError {
status: 429,
message: "try later".into(),
reason: None,
})
.context("obtain native coordination grant");
let unavailable = anyhow::Error::new(ControlPlaneHttpError {
status: 503,
message: "provider unavailable".into(),
reason: None,
})
.context("obtain native coordination grant");
assert!(!is_retryable_gateway_error(&revoked));
assert!(is_retryable_gateway_error(&throttled));
assert!(is_retryable_gateway_error(&unavailable));
}
#[test]
fn lease_fallback_error_owns_terminal_reconnect_classification() {
let permanent_lease_error = gateway_connect_error("missing operation price", false);
let transient_fallback_error = gateway_connect_error("control plane timeout", true);
let terminal =
terminal_refresh_error(permanent_lease_error, Some(transient_fallback_error));
assert!(is_retryable_gateway_error(&terminal));
assert_eq!(terminal.to_string(), "control plane timeout");
}
#[tokio::test]
async fn native_keepalive_is_a_websocket_control_frame() {
let (mut sender, mut receiver) = futures::channel::mpsc::channel(1);
send_socket_keepalive(&mut sender).await.unwrap();
assert!(
matches!(receiver.next().await, Some(Message::Ping(payload)) if payload.is_empty())
);
}
#[test]
fn only_local_device_delete_terminates_the_gateway_presence_owner() {
let (local_reply, _local_result) = oneshot::channel();
let local = Command::Delete {
user_id: "user-1".into(),
device_id: "local-device".into(),
reply: local_reply,
};
let (remote_reply, _remote_result) = oneshot::channel();
let remote = Command::Delete {
user_id: "user-1".into(),
device_id: "remote-device".into(),
reply: remote_reply,
};
assert!(command_deletes_device(&local, "local-device"));
assert!(!command_deletes_device(&remote, "local-device"));
}
#[tokio::test]
async fn healthy_lease_response_advances_the_local_device_projection() {
let provider = Arc::new(RecordingGrantProvider::default());
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-1".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "user".to_string(),
id: "principal-1".to_string(),
},
architecture: None,
room_delivery: None,
grant_provider: provider,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
adapter.shared.devices.write().await.insert(
"device-1".to_string(),
Device {
app_tag: Some("app_native_test".to_string()),
device_id: "device-1".to_string(),
user_id: Some("principal-1".to_string()),
device_name: "Native".to_string(),
platform_type: Some("desktop".to_string()),
capabilities: None,
session_id: Some("runtime-1".to_string()),
node_id: Some("00".repeat(32)),
tag: None,
kind: None,
metadata: None,
online: true,
ticket: Some("ticket-1".to_string()),
last_seen_at: None,
expires_at: Some(serde_json::json!(1_000_u64)),
created_at: None,
updated_at: None,
excluded_peers: Vec::new(),
},
);
update_local_device_expiry(&adapter, 2_000).await;
let devices = adapter.shared.devices.read().await;
assert_eq!(
devices
.get("device-1")
.and_then(|device| device.expires_at.clone()),
Some(serde_json::json!(2_000_u64)),
);
}
#[tokio::test]
async fn native_membership_pages_publish_only_after_the_complete_snapshot() {
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-local".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "room".to_string(),
id: "room-paged".to_string(),
},
architecture: Some(RoomArchitectureMode::Managed),
room_delivery: Some(RoomDelivery::Reliable),
grant_provider: Arc::new(RecordingGrantProvider::default()),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
let members = (0..201)
.map(|index| GatewayMember {
user_id: Some(format!("user-{index}")),
device_id: format!("device-{index:03}"),
device_name: format!("Device {index}"),
platform_type: "test".to_string(),
metadata: None,
capabilities: None,
online: true,
updated_at_ms: 1_000,
expires_at_ms: 60_000,
})
.collect::<Vec<_>>();
for page_index in 0..3 {
handle_frame(
&adapter,
ServerFrame::MembershipPage {
snapshot_id: "connection:paged-room".to_string(),
page_index,
page_count: 3,
member_count: members.len(),
members: members[page_index * 100..((page_index + 1) * 100).min(members.len())]
.to_vec(),
},
)
.await
.expect("membership page");
if page_index < 2 {
assert!(adapter.shared.members.read().await.is_empty());
}
}
assert_eq!(adapter.shared.members.read().await.len(), 201);
assert!(adapter
.shared
.membership_pages
.lock()
.expect("membership page lock")
.is_none());
}
#[tokio::test]
async fn v2_topology_projects_room_architecture_and_selected_routes() {
let provider = Arc::new(RecordingGrantProvider::default());
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".to_string(),
app_tag: "app_native_test".to_string(),
device_id: "device-local".to_string(),
platform_type: "desktop".to_string(),
avenue: NativeCoordinationAvenue {
kind: "room".to_string(),
id: "room-adaptive".to_string(),
},
architecture: Some(RoomArchitectureMode::Auto),
room_delivery: Some(RoomDelivery::LatestState),
grant_provider: provider,
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
})
.expect("adapter");
let now = now_ms();
adapter.shared.members.write().await.insert(
"device-peer".to_string(),
GatewayMember {
user_id: None,
device_id: "device-peer".to_string(),
device_name: "Peer".to_string(),
platform_type: "desktop".to_string(),
metadata: None,
capabilities: None,
online: true,
updated_at_ms: now as i64,
expires_at_ms: now.saturating_add(60_000) as i64,
},
);
*adapter.shared.active_grant.write().await =
Some(("grant-current".to_string(), now.saturating_add(60_000)));
let current_lease = GatewayTopologyLease {
schema_version: 2,
topology_revision: 1,
previous_topology_revision: None,
avenue: adapter.avenue.clone(),
grant_jti: "grant-current".to_string(),
expires_at_ms: now.saturating_add(30_000),
architecture: GatewayArchitectureLease {
policy_version: "room-architecture-v1".to_string(),
epoch: 1,
previous_architecture_epoch: None,
requested_mode: RoomArchitectureMode::Auto,
effective_mode: EffectiveRoomArchitecture::Sparse,
phase: RoomArchitecturePhase::Settled,
reason: RoomArchitectureReason::Size,
held_credits_microusd: 25_000,
quote_expires_at_ms: now.saturating_add(30_000),
reservation_id: Some("reservation:test".to_string()),
},
group_encryption: None,
active: vec![GatewayDevice {
user_id: None,
device_id: "device-peer".to_string(),
runtime_instance_id: "runtime-peer".to_string(),
node_id: "node-peer".to_string(),
device_name: "Peer".to_string(),
platform_type: "desktop".to_string(),
ticket: "ticket-peer".to_string(),
metadata: None,
capabilities: None,
excluded_peers: Vec::new(),
online: true,
updated_at_ms: now as i64,
expires_at_ms: now.saturating_add(60_000) as i64,
}],
backups: Vec::new(),
};
let mut wrong_grant = current_lease.clone();
wrong_grant.grant_jti = "grant-forged".to_string();
assert!(accept_topology_lease(&adapter, wrong_grant).await.is_err());
assert_eq!(adapter.room_architecture(), None);
accept_topology_lease(&adapter, current_lease.clone())
.await
.expect("v2 topology lease");
assert_eq!(
adapter.room_architecture(),
Some(RoomArchitectureSnapshot {
requested: RoomArchitectureMode::Auto,
effective: EffectiveRoomArchitecture::Sparse,
epoch: 1,
phase: RoomArchitecturePhase::Settled,
reason: RoomArchitectureReason::Size,
held_credits_usd: 0.025,
quote_expires_at_ms: now.saturating_add(30_000),
}),
);
assert_eq!(
adapter
.shared
.devices
.read()
.await
.get("device-peer")
.and_then(|device| device.node_id.as_deref()),
Some("node-peer"),
);
accept_topology_lease(&adapter, current_lease.clone())
.await
.expect("exact replay is stale and idempotent");
assert_eq!(
adapter.room_architecture().map(|snapshot| snapshot.epoch),
Some(1),
);
let mut forged_replay = current_lease;
forged_replay.grant_jti = "grant-forged".to_string();
assert!(accept_topology_lease(&adapter, forged_replay)
.await
.is_err());
assert_eq!(
adapter.room_architecture().map(|snapshot| snapshot.epoch),
Some(1),
);
let wrong_predecessor = GatewayTopologyLease {
schema_version: 2,
topology_revision: 2,
previous_topology_revision: Some(99),
avenue: adapter.avenue.clone(),
grant_jti: "grant-current".to_string(),
expires_at_ms: now.saturating_add(30_000),
architecture: GatewayArchitectureLease {
policy_version: "room-architecture-v1".to_string(),
epoch: 2,
previous_architecture_epoch: Some(1),
requested_mode: RoomArchitectureMode::Auto,
effective_mode: EffectiveRoomArchitecture::Managed,
phase: RoomArchitecturePhase::Preparing,
reason: RoomArchitectureReason::Size,
held_credits_microusd: 25_000,
quote_expires_at_ms: now.saturating_add(30_000),
reservation_id: Some("reservation:managed".to_string()),
},
group_encryption: None,
active: Vec::new(),
backups: Vec::new(),
};
assert!(accept_topology_lease(&adapter, wrong_predecessor)
.await
.is_err());
assert_eq!(
adapter.room_architecture().map(|snapshot| snapshot.epoch),
Some(1),
);
accept_topology_lease(
&adapter,
GatewayTopologyLease {
schema_version: 2,
topology_revision: 2,
previous_topology_revision: Some(1),
avenue: adapter.avenue.clone(),
grant_jti: "grant-current".to_string(),
expires_at_ms: now.saturating_add(30_000),
architecture: GatewayArchitectureLease {
policy_version: "room-architecture-v1".to_string(),
epoch: 2,
previous_architecture_epoch: Some(1),
requested_mode: RoomArchitectureMode::Auto,
effective_mode: EffectiveRoomArchitecture::Managed,
phase: RoomArchitecturePhase::Preparing,
reason: RoomArchitectureReason::Size,
held_credits_microusd: 25_000,
quote_expires_at_ms: now.saturating_add(30_000),
reservation_id: Some("reservation:managed".to_string()),
},
group_encryption: None,
active: vec![GatewayDevice {
user_id: None,
device_id: "device-peer".to_string(),
runtime_instance_id: "runtime-peer".to_string(),
node_id: "node-peer".to_string(),
device_name: "Peer".to_string(),
platform_type: "desktop".to_string(),
ticket: "ticket-peer".to_string(),
metadata: None,
capabilities: None,
excluded_peers: Vec::new(),
online: true,
updated_at_ms: now as i64,
expires_at_ms: now.saturating_add(60_000) as i64,
}],
backups: Vec::new(),
},
)
.await
.expect("preparing managed lease keeps sparse routes without a group lease");
assert_eq!(
adapter
.room_architecture()
.map(|snapshot| (snapshot.effective, snapshot.phase)),
Some((
EffectiveRoomArchitecture::Managed,
RoomArchitecturePhase::Preparing,
)),
);
assert!(adapter
.shared
.devices
.read()
.await
.contains_key("device-peer"));
}
#[tokio::test]
async fn command_completion_preserves_permanent_error_classification() {
let (reply, result) = oneshot::channel();
let actor_error = complete_command::<()>(
reply,
Err(gateway_connect_error("invalid payload", false)),
"presence publication",
)
.expect_err("permanent operation must fail");
assert!(!is_retryable_gateway_error(&actor_error));
assert_eq!(
result
.await
.expect("caller receives reply")
.expect_err("caller receives operation failure")
.to_string(),
"invalid payload"
);
}
#[test]
fn optional_device_fields_are_omitted_instead_of_null() {
let value = serde_json::to_value(OutboundGatewayDevice {
device_id: "device:test".into(),
runtime_instance_id: "runtime:test".into(),
node_id: "00".repeat(32),
device_name: "Test".into(),
platform_type: "macos".into(),
ticket: "ticket".into(),
metadata: None,
capabilities: None,
excluded_peers: Vec::new(),
online: true,
})
.unwrap();
assert!(value.get("metadata").is_none());
assert!(value.get("capabilities").is_none());
}
}