#[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 AUTHORITY_TICKET_RENEWAL_SKEW_MS: u64 = 60_000;
const SOCKET_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30);
#[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, Clone)]
pub struct GatewayConnectError {
message: String,
retryable: bool,
code: Option<String>,
pub service_error: Option<crate::service_errors::ServiceError>,
}
impl fmt::Display for GatewayConnectError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.message)
}
}
impl StdError for GatewayConnectError {}
fn copy_gateway_failure(error: &anyhow::Error) -> anyhow::Error {
let detail = format!("{error:#}");
if let Some(denial) = error.downcast_ref::<GatewayConnectError>() {
return anyhow::Error::new(denial.clone()).context(detail);
}
if let Some(denial) = error.downcast_ref::<ControlPlaneHttpError>() {
return anyhow::Error::new(ControlPlaneHttpError {
status: denial.status,
message: denial.message.clone(),
reason: denial.reason.clone(),
service_error: denial.service_error.clone(),
}).context(detail);
}
anyhow!(detail)
}
fn gateway_connect_error(message: impl Into<String>, retryable: bool) -> anyhow::Error {
anyhow!(GatewayConnectError {
message: message.into(),
retryable,
code: None,
service_error: None,
})
}
fn gateway_server_error(code: String, message: String, retryable: bool) -> anyhow::Error {
gateway_server_error_metadata(code, message, retryable, Default::default())
}
fn gateway_server_error_metadata(
code: String,
message: String,
retryable: bool,
mut metadata: serde_json::Map<String, serde_json::Value>,
) -> anyhow::Error {
metadata.insert("code".into(), serde_json::Value::String(code.clone()));
metadata.insert("retryable".into(), serde_json::Value::Bool(retryable));
let service_error = crate::service_errors::ServiceError::from_value(&metadata.into());
anyhow!(GatewayConnectError {
message: format!("coordination gateway {code}: {message}"),
retryable,
code: Some(code),
service_error,
})
}
fn is_retryable_gateway_error(error: &anyhow::Error) -> bool {
if is_gateway_revocation(error) {
return false;
}
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 gateway_service_error(error: &anyhow::Error) -> Option<&crate::service_errors::ServiceError> {
crate::service_errors::ServiceError::from_error(error)
}
fn gateway_service_retry_delay(error: &anyhow::Error) -> Duration {
let Some(service) = gateway_service_error(error) else {
return Duration::ZERO;
};
Duration::from_millis(
service.retry_after_ms.unwrap_or(0).max(
service.reset_at.unwrap_or(0).saturating_sub(now_ms()),
),
)
}
fn is_gateway_revocation(error: &anyhow::Error) -> bool {
error
.downcast_ref::<GatewayConnectError>()
.is_some_and(|error| error.code.as_deref() == Some("credential-revoked"))
|| error
.downcast_ref::<ControlPlaneHttpError>()
.is_some_and(|error| error.reason.as_deref() == Some("device-certificate-revoked"))
}
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, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeServiceErrorObservation {
pub avenue: NativeCoordinationAvenue,
pub runtime_instance_id: String,
pub service_error: crate::service_errors::ServiceError,
}
#[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(crate) async fn compose_authority_client(
&self,
builder: crate::client::ClientBuilder,
) -> Result<Arc<crate::Client>> {
if self.is_closed()
|| self.kind != NativeCapabilityKind::Room
|| self.signaling.architecture != Some(RoomArchitectureMode::Authority)
{
bail!("native authority composition requires an open authority room");
}
let mut binding = self.signaling.authority_client.lock().await;
if self.is_closed() {
bail!("native authority is closed");
}
if binding.is_some() {
bail!("native authority already has a peer-session owner");
}
let client = Arc::new(builder.signaling_backend(self.signaling()).build());
if client.app_tag() != self.signaling.app_tag {
bail!("native authority client app identity does not match its capability");
}
let scope = format!("v2:room:{}", self.id);
client
.session_token_registry
.require_scope_peer_admission(&scope)
.map_err(anyhow::Error::msg)?;
client.ensure_default_admission_gate(&scope);
client
.start_external_auto_connect(
format!("{}:native-root", self.signaling.app_tag),
self.signaling.device_id.clone(),
)
.await?;
*binding = Some(AuthorityPeerBinding {
client: Arc::downgrade(&client),
scope,
revision: 0,
admission_revision: 0,
last_payload: String::new(),
retiring_tickets: Vec::new(),
closed: false,
});
Ok(client)
}
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())
}
pub fn subscribe_service_errors(
&self,
) -> Result<broadcast::Receiver<NativeServiceErrorObservation>> {
if self.is_closed() {
bail!("native capability is closed");
}
Ok(self.signaling.subscribe_service_errors())
}
#[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>,
#[serde(default)]
admission_peers: Option<Vec<GatewayAdmissionPeer>>,
active: Vec<GatewayDevice>,
backups: Vec<GatewayDevice>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GatewayAdmissionPeer {
device_id: String,
node_id: String,
}
#[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(flatten)]
metadata: serde_json::Map<String, serde_json::Value>,
},
#[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, bool)>>,
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>>,
service_errors: broadcast::Sender<NativeServiceErrorObservation>,
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>,
}
struct AuthorityPeerBinding {
client: std::sync::Weak<crate::Client>,
scope: String,
revision: u64,
admission_revision: u64,
last_payload: String,
retiring_tickets: Vec<(String, u64)>,
closed: bool,
}
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>,
authority_client: tokio::sync::Mutex<Option<AuthorityPeerBinding>>,
#[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 (service_errors, _) = 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,
service_errors,
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,
authority_client: tokio::sync::Mutex::new(None),
#[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) {
self.stop_authority_client().await;
let _ = self.commands.send(Command::Stop).await;
}
async fn authority_ticket_maintenance_delay(&self) -> Duration {
let binding = self.authority_client.lock().await;
let Some(binding) = binding.as_ref().filter(|binding| !binding.closed && binding.client.strong_count() > 0) else {
return Duration::from_secs(365 * 24 * 60 * 60);
};
let desired = self.shared.desired.read().await;
let renewal = desired.as_ref().map(|desired| {
let (_, suffix) = crate::session_token::split_ticket(&desired.ticket);
suffix
.and_then(crate::session_token::decode_token_payload)
.and_then(|payload| payload.expires_at_ms)
.map(|expires| expires.saturating_sub(AUTHORITY_TICKET_RENEWAL_SKEW_MS))
.unwrap_or(0)
});
let deadline = binding
.retiring_tickets
.iter()
.map(|(_, expiry)| *expiry)
.chain(renewal)
.min();
deadline.map_or(Duration::from_secs(365 * 24 * 60 * 60), |deadline| {
Duration::from_millis(deadline.saturating_sub(now_ms()).max(1))
})
}
async fn maintain_authority_ticket(&self) -> Result<bool> {
self.maintain_authority_publication(None).await
}
async fn maintain_authority_publication(&self, pending: Option<&mut DesiredPresence>) -> Result<bool> {
self.expire_retiring_authority_tickets().await;
let mut binding = self.authority_client.lock().await;
let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
return Ok(false);
};
let Some(client) = binding.client.upgrade() else {
return Ok(false);
};
let now = now_ms();
let mut desired = self.shared.desired.write().await;
let Some(desired) = desired.as_mut() else {
return Ok(false);
};
let (ticket, suffix) = crate::session_token::split_ticket(&desired.ticket);
let payload = suffix
.and_then(crate::session_token::decode_token_payload)
.ok_or_else(|| anyhow!("authority presence requires a scoped expiring ticket"))?;
if payload.scope.as_str() != binding.scope
|| payload.ticket_hash.as_deref()
!= Some(crate::session_token::endpoint_ticket_hash(ticket).as_str())
{
bail!("authority presence ticket scope or endpoint binding is invalid");
}
let expires = payload
.expires_at_ms
.ok_or_else(|| anyhow!("authority ticket must expire"))?;
if expires.saturating_sub(now) > AUTHORITY_TICKET_RENEWAL_SKEW_MS {
return Ok(false);
}
if binding.retiring_tickets.len() >= 2 {
bail!("authority ticket overlap capacity exhausted");
}
let replacement = client
.endpoint_ticket_with_token(&binding.scope, payload.max_connections)
.await?;
if let Some(pending) = pending.filter(|pending| *pending == desired) {
pending.ticket = replacement.clone();
}
desired.ticket = replacement;
if expires <= now {
client.revoke_session_token(&payload.token);
} else {
binding.retiring_tickets.push((payload.token, expires));
}
Ok(true)
}
async fn expire_retiring_authority_tickets(&self) -> Duration {
let mut binding = self.authority_client.lock().await;
let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
return Duration::from_secs(365 * 24 * 60 * 60);
};
let Some(client) = binding.client.upgrade() else {
return Duration::from_secs(365 * 24 * 60 * 60);
};
let now = now_ms();
binding.retiring_tickets.retain(|(token, expiry)| {
if *expiry > now { return true; }
client.revoke_session_token(token);
false
});
binding.retiring_tickets.iter().map(|(_, expiry)| *expiry).min()
.map_or(Duration::from_secs(365 * 24 * 60 * 60), |expiry| {
Duration::from_millis(expiry.saturating_sub(now).max(1))
})
}
async fn forward_authority_peers(&self) -> Result<()> {
let mut binding = self.authority_client.lock().await;
let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
return Ok(());
};
let Some(client) = binding.client.upgrade() else {
return Ok(());
};
let mut peers = Vec::new();
if self.room_architecture().is_some_and(|snapshot| {
snapshot.effective == EffectiveRoomArchitecture::Authority
&& snapshot.phase == RoomArchitecturePhase::Settled
}) {
let members = self.shared.members.read().await;
let routes = self.shared.topology_routes.read().await;
for (device, active) in routes.values() {
if !active
|| !device.online
|| device.device_id == self.device_id
|| !members
.get(&device.device_id)
.is_some_and(|member| member.online)
{
continue;
}
peers.push(serde_json::json!({
"deviceId": device.device_id,
"nodeId": device.node_id,
"ticket": device.ticket,
"online": true,
"sessionId": device.runtime_instance_id,
"excludedPeers": device.excluded_peers,
}));
}
}
peers.sort_by(|left, right| left["deviceId"].as_str().cmp(&right["deviceId"].as_str()));
let payload = serde_json::to_string(&peers)?;
if payload != binding.last_payload {
let revision = binding
.revision
.checked_add(1)
.ok_or_else(|| anyhow!("native authority input revision exhausted"))?;
if !client
.submit_external_desired_peers(revision, &payload)
.await?
{
bail!("native authority peer input rejected by its runtime owner");
}
binding.revision = revision;
binding.last_payload = payload;
}
Ok(())
}
async fn handle_gateway_failure(&self, error: &anyhow::Error) -> bool {
if let Some(service_error) = gateway_service_error(error) {
let _ = self.shared.service_errors.send(NativeServiceErrorObservation {
avenue: self.avenue.clone(),
runtime_instance_id: self.runtime_instance_id.clone(),
service_error: service_error.clone(),
});
}
if is_gateway_revocation(error) {
self.stop_authority_client().await;
}
is_retryable_gateway_error(error)
}
fn subscribe_service_errors(&self) -> broadcast::Receiver<NativeServiceErrorObservation> {
self.shared.service_errors.subscribe()
}
async fn stop_authority_client(&self) {
let mut binding = self.authority_client.lock().await;
let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
return;
};
binding.closed = true;
if let Some(client) = binding.client.upgrade() {
binding.admission_revision += 1;
if let Err(error) = client.session_token_registry.update_scope_peer_admission(
&binding.scope,
binding.admission_revision,
0,
Vec::new(),
) {
eprintln!("[openrtc][authority-close] admission withdrawal failed: {error}");
}
client.register_session_token(
crate::session_token::generate_token(),
"authority-closed".to_string(),
0,
);
let revoked = client.begin_revoke_tokens_by_scope(&binding.scope);
if let Some(revision) = binding.revision.checked_add(1) {
if let Err(error) = client.submit_external_desired_peers(revision, "[]").await {
eprintln!("[openrtc][authority-close] peer withdrawal failed: {error:#}");
}
}
client.stop_external_auto_connect().await;
client
.finish_revoke_tokens_by_scope(&binding.scope, &revoked)
.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 = {
let mut routes = adapter.shared.topology_routes.write().await;
routes.retain(|device_id, _| members.contains_key(device_id));
routes.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, _active)| 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 authority_return_retry_available = false;
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, anyhow::Error)> = None;
let mut keepalive = tokio::time::interval(SOCKET_KEEPALIVE_INTERVAL);
keepalive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
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(copy_gateway_failure(message)));
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));
if retry_at.is_none() {
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, pending_publish.as_mut().map(|(desired, _)| desired)).await;
match result {
Ok(connected) => {
reconnect_attempt = 0;
authority_return_retry_available = true;
circuit_failure = None;
*adapter.shared.applied_presence.write().await = adapter.shared.desired.read().await.clone();
socket = Some(connected);
if let Some((published, reply)) = pending_publish.take() {
if Some(&published) == adapter.shared.applied_presence.read().await.as_ref() {
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 = adapter.handle_gateway_failure(&error).await;
eprintln!(
"[openrtc][coordination-gateway][connect-retry] attempt={}/{} retryable={} error={}",
reconnect_attempt,
MAX_RECONNECT_ATTEMPTS,
retryable,
error
);
let delay = gateway_reconnect_delay(
reconnect_attempt, &error, &mut authority_return_retry_available,
);
if let Some(delay) = delay {
retry_at = Some(tokio::time::Instant::now() + delay);
} else {
retry_at = None;
let failure = error.context(format!(
"native coordination gateway unavailable after {} attempts",
reconnect_attempt
));
if let Some((_, reply)) = pending_publish.take() {
let _ = reply.send(Err(copy_gateway_failure(&failure)));
}
circuit_failure = Some((desired, failure));
}
}
}
}
_ = tokio::time::sleep(adapter.authority_ticket_maintenance_delay().await) => {
if let Err(error) = adapter.maintain_authority_publication(
pending_publish.as_mut().map(|(desired, _)| desired),
).await {
eprintln!("[openrtc][coordination-gateway][offline-ticket-failed] error={error:#}");
adapter.stop_authority_client().await;
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, copy_gateway_failure(&error)));
}
if let Some((_, reply)) = pending_publish.take() {
let _ = reply.send(Err(error));
}
}
}
}
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(copy_gateway_failure(message)));
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 adapter.handle_gateway_failure(&error).await {
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
} else {
eprintln!(
"[openrtc][coordination-gateway][operation-rejected] retryable=false error={}",
error
);
if is_gateway_revocation(&error) {
let _ = active.socket.close(None).await;
*adapter.shared.applied_presence.write().await = None;
socket = None;
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error));
}
continue;
}
if publishing {
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error));
}
}
}
}
}
_ = tokio::time::sleep(adapter.authority_ticket_maintenance_delay().await) => {
let renewal = async {
if adapter.maintain_authority_ticket().await? {
let desired = adapter.shared.desired.read().await.clone()
.ok_or_else(|| anyhow!("authority presence was withdrawn"))?;
refresh_presence_authentication(&adapter, active, desired).await?;
}
Ok::<_, anyhow::Error>(())
}.await;
if let Err(error) = renewal {
let retryable = adapter.handle_gateway_failure(&error).await;
eprintln!("[openrtc][coordination-gateway][authority-ticket-failed] retryable={} error={error:#}", retryable);
*adapter.shared.applied_presence.write().await = None;
let _ = active.socket.close(None).await;
socket = None;
reconnect_attempt = 0;
if retryable {
retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
} else {
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error));
}
}
}
}
_ = 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
&& !is_gateway_revocation(&error)
&& gateway_service_error(&error).is_none() {
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
};
let retryable = adapter.handle_gateway_failure(&terminal_error).await;
eprintln!(
"[openrtc][coordination-gateway][auth-refresh-failed] retryable={} error={}",
retryable, terminal_error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
if retryable {
retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&terminal_error));
} else {
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, terminal_error));
}
}
}
}
}
_ = keepalive.tick() => {
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 {
let retryable = adapter.handle_gateway_failure(&error).await;
eprintln!(
"[openrtc][coordination-gateway][socket-frame-failed] retryable={} error={:#}",
retryable,
error,
);
*adapter.shared.applied_presence.write().await = None;
socket = None;
reconnect_attempt = 0;
if retryable {
retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
} else {
retry_at = None;
if let Some(desired) = adapter.shared.desired.read().await.clone() {
circuit_failure = Some((desired, error));
}
}
}
}
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());
},
}
}
}
}
}
adapter.stop_authority_client().await;
}
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, pending: Option<&mut DesiredPresence>) -> Result<ConnectedGateway> {
adapter.maintain_authority_publication(pending).await?;
let desired = adapter
.shared
.desired
.read()
.await
.clone()
.ok_or_else(|| anyhow!("native coordination presence is not configured"))?;
let work = connect_with_desired(adapter, &desired, "presence");
tokio::pin!(work);
let mut expiry_delay = adapter.expire_retiring_authority_tickets().await;
loop {
tokio::select! {
result = &mut work => return result,
delay = async {
tokio::time::sleep(expiry_delay).await;
adapter.expire_retiring_authority_tickets().await
} => expiry_delay = delay,
}
}
}
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,
"socketLiveness": "ping-v1",
"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,
metadata,
..
} => {
return Err(gateway_server_error_metadata(code, message, retryable, metadata));
}
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"))?;
refresh_authentication_for_presence(adapter, socket, current_grant, &desired).await
}
async fn refresh_authentication_for_presence(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
current_grant: Option<String>,
desired: &DesiredPresence,
) -> Result<(u64, Option<String>)> {
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,
metadata,
..
} => {
return Err(gateway_server_error_metadata(code, message, retryable, metadata));
}
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_presence_authentication(
adapter: &NativeCoordinationGatewaySignaling,
active: &mut ConnectedGateway,
desired: DesiredPresence,
) -> Result<()> {
let (expires_at_ms, token) = refresh_authentication_for_presence(
adapter,
&mut active.socket,
Some(active.credential_token.clone()),
&desired,
)
.await?;
active.credential_expires_at_ms = expires_at_ms;
if let Some(token) = token {
active.credential_token = token;
}
execute_operation_for_presence(
adapter,
&mut active.socket,
presence_frame(adapter, &desired).await,
Some(&desired),
)
.await?;
*adapter.shared.applied_presence.write().await = Some(desired);
Ok(())
}
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 => {
let mut grant = adapter.shared.active_grant.write().await;
let (_, current_expiry) = grant.as_mut().ok_or_else(|| {
gateway_connect_error("lease refresh has no current grant", false)
})?;
if expires_at_ms <= now_ms() {
return Err(gateway_connect_error(
"lease refresh is already expired",
false,
));
}
*current_expiry = expires_at_ms;
drop(grant);
update_local_device_expiry(adapter, expires_at_ms).await;
return Ok((expires_at_ms, None));
}
ServerFrame::Error {
code,
message,
retryable,
metadata,
..
} => {
return Err(gateway_server_error_metadata(code, message, retryable, metadata));
}
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 message = format!("{error:#}");
let _ = reply.send(Err(copy_gateway_failure(&error)));
Err(error.context(format!("{operation} failed: {message}")))
}
}
}
async fn execute_operation(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
frame: serde_json::Value,
) -> Result<()> {
execute_operation_for_presence(adapter, socket, frame, None).await
}
async fn execute_operation_for_presence(
adapter: &NativeCoordinationGatewaySignaling,
socket: &mut GatewaySocket,
frame: serde_json::Value,
bound_presence: Option<&DesiredPresence>,
) -> 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,
metadata,
} if key == idempotency_key => {
if operation_requires_budget_renewal(&code, retryable) {
return Ok(true);
}
return Err(gateway_server_error_metadata(
code,
format!("{message} operation={operation}"),
retryable,
metadata,
));
}
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}");
}
let renewal = match bound_presence {
Some(desired) => {
refresh_authentication_for_presence(adapter, socket, None, desired).await
}
None => refresh_authentication(adapter, socket, None).await,
};
renewal
.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_peer_admission = architecture.effective_mode != EffectiveRoomArchitecture::Authority
|| lease.admission_peers.as_ref().is_some_and(|peers| {
peers.len() <= 50
&& peers
.iter()
.map(|peer| &peer.device_id)
.collect::<std::collections::HashSet<_>>()
.len()
== peers.len()
&& peers.iter().all(|peer| {
peer.device_id != adapter.device_id
&& members
.get(&peer.device_id)
.is_some_and(|member| member.online)
&& peer.node_id.parse::<iroh::EndpointId>().is_ok()
})
});
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
|| !valid_peer_admission
{
return Err(invalid());
}
drop(members);
if adapter.authority_client.lock().await.is_some() {
let scope = format!("v2:room:{}", adapter.avenue.id);
if routes.iter().any(|route| {
let (ticket, suffix) = crate::session_token::split_ticket(&route.ticket);
!suffix
.and_then(|suffix| crate::session_token::decode_payload(ticket, suffix))
.is_some_and(|payload| payload.scope.as_str() == scope)
}) {
return Err(gateway_connect_error(
"native authority route lacks room-scoped admission",
false,
));
}
}
if architecture.effective_mode == EffectiveRoomArchitecture::Authority {
let mut binding = adapter.authority_client.lock().await;
if let Some(binding) = binding.as_mut() {
let client = binding
.client
.upgrade()
.filter(|_| !binding.closed)
.ok_or_else(invalid)?;
binding.admission_revision = binding
.admission_revision
.checked_add(1)
.ok_or_else(invalid)?;
client
.session_token_registry
.update_scope_peer_admission(
&binding.scope,
binding.admission_revision,
lease.expires_at_ms,
lease
.admission_peers
.as_ref()
.ok_or_else(invalid)?
.iter()
.map(|peer| peer.node_id.clone())
.collect(),
)
.map_err(anyhow::Error::msg)?;
}
}
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()
.map(|route| (route, true))
.chain(lease.backups.into_iter().map(|route| (route, false)))
.map(|(route, active)| (route.device_id.clone(), (route, active)))
.collect();
*adapter.shared.topology_revision.write().await = lease.topology_revision;
rebuild_sparse_devices(adapter).await;
adapter.forward_authority_peers().await
}
fn decode_frame(message: Message) -> Result<Option<ServerFrame>> {
match message {
Message::Text(text) if text == "pong" => Ok(None),
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(Some(frame)) if u16::from(frame.code) == 4403 => Err(gateway_server_error(
"credential-revoked".into(),
"gateway authorization withdrawn".into(),
false,
)),
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,
metadata,
..
} => return Err(gateway_server_error_metadata(code, message, retryable, metadata)),
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 => {}
}
adapter.forward_authority_peers().await
}
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 gateway_reconnect_delay(
attempt: u8,
error: &anyhow::Error,
authority_return_retry_available: &mut bool,
) -> Option<Duration> {
if !is_retryable_gateway_error(error) || attempt >= MAX_RECONNECT_ATTEMPTS {
return None;
}
if *authority_return_retry_available
&& error
.downcast_ref::<GatewayConnectError>()
.is_some_and(|error| error.code.as_deref() == Some("room-authority-unavailable"))
{
*authority_return_retry_available = false;
return Some(reconnect_delay(1));
}
Some(reconnect_delay(attempt).max(gateway_service_retry_delay(error)))
}
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::Text("ping".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>>,
gateway_url: Option<String>,
token: Option<String>,
blocked: Option<Arc<tokio::sync::Notify>>,
}
#[async_trait]
impl NativeGatewayGrantProvider for RecordingGrantProvider {
async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
self.requests.lock().expect("requests").push(request);
if let Some(blocked) = &self.blocked { blocked.notified().await; }
Ok(NativeGatewayGrant {
protocol_version: 2,
gateway_url: self
.gateway_url
.clone()
.unwrap_or_else(|| "https://gateway.example.test".into()),
route_key: "route-1".to_string(),
token: self.token.clone().unwrap_or_else(|| "grant-2".into()),
expires_at_ms: now_ms() + 3_600_000,
})
}
}
#[tokio::test]
async fn socket_lease_refresh_updates_only_acknowledged_grant_expiry() {
for outcome in ["accepted", "wrong-request", "denied", "expired"] {
let provider = Arc::new(RecordingGrantProvider::default());
let handle = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
"lease-client",
"native",
provider.clone(),
)
.unwrap()
.join_room_with_architecture("lease-room", RoomArchitectureMode::Authority)
.unwrap();
let adapter = &handle.signaling;
let old_expiry = now_ms() + 30_000;
let renewed_expiry = if outcome == "expired" {
1
} else {
old_expiry + 60_000
};
*adapter.shared.active_grant.write().await =
Some(("same-socket-grant".into(), old_expiry));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (tcp, _) = listener.accept().await.unwrap();
let mut socket = tokio_tungstenite::accept_async(tcp).await.unwrap();
let request: serde_json::Value =
serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
.unwrap();
assert_eq!(request["type"], "lease.refresh");
let response = if outcome == "denied" {
serde_json::json!({"type":"error", "code":"credential-revoked", "message":"local denial", "retryable":false})
} else {
serde_json::json!({"type":"lease.refreshed", "expiresAtMs":renewed_expiry,
"idempotencyKey":if outcome == "wrong-request" { serde_json::json!("other-request") } else { request["idempotencyKey"].clone() }})
};
socket
.send(Message::Text(response.to_string().into()))
.await
.unwrap();
socket.close(None).await.unwrap();
});
let (mut socket, _) = tokio_tungstenite::connect_async(format!("ws://{address}"))
.await
.unwrap();
let result =
tokio::time::timeout(Duration::from_secs(3), refresh_lease(adapter, &mut socket))
.await
.unwrap();
let accepted = outcome == "accepted";
assert_eq!(result.is_ok(), accepted, "{outcome}: {result:?}");
assert_eq!(*adapter.shared.active_grant.read().await,
Some(("same-socket-grant".into(), if accepted { renewed_expiry } else { old_expiry })),
"{outcome}: topology validation must use only the acknowledged socket authorization");
assert!(
provider.requests.lock().unwrap().is_empty(),
"lease renewal must not mint a grant"
);
tokio::time::timeout(Duration::from_secs(3), server)
.await
.unwrap()
.unwrap();
handle.close().await;
}
}
#[tokio::test]
async fn authority_ticket_refresh_keeps_socket_and_orders_authorization_before_presence() {
for outcome in [
"accepted",
"budget-replay",
"auth-denied",
"presence-denied",
"auth-eof",
] {
let accepted = matches!(outcome, "accepted" | "budget-replay");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let endpoint = format!("http://{}", listener.local_addr().unwrap());
let token = format!(
"local.{}.test",
base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(br#"{"jti":"renewed-local-grant"}"#)
);
let provider = Arc::new(RecordingGrantProvider {
gateway_url: Some(endpoint.clone()),
token: Some(token.clone()),
..Default::default()
});
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: endpoint.clone(),
app_tag: "app_native_test".into(),
device_id: "service-device".into(),
platform_type: "native".into(),
avenue: NativeCoordinationAvenue {
kind: "room".into(),
id: "local-room".into(),
},
architecture: Some(RoomArchitectureMode::Authority),
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,
})
.unwrap();
let old = DesiredPresence {
user_id: "local-room".into(),
local_node_id: "local-node".into(),
ticket: "old-local-ticket".into(),
device_name: "Service".into(),
metadata: None,
ttl_ms: 900_000,
online: true,
};
let desired = DesiredPresence {
ticket: "renewed-local-ticket".into(),
..old.clone()
};
*adapter.shared.desired.write().await = Some(desired.clone());
*adapter.shared.applied_presence.write().await = Some(old.clone());
let server_adapter = adapter.clone();
let server = tokio::spawn(async move {
let (tcp, _) = listener.accept().await.unwrap();
let mut socket = tokio_tungstenite::accept_async(tcp).await.unwrap();
let frame: serde_json::Value =
serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
.unwrap();
assert_eq!(frame["type"], "auth.refresh");
assert_eq!(
server_adapter
.shared
.applied_presence
.read()
.await
.as_ref()
.unwrap()
.ticket,
"old-local-ticket"
);
if outcome == "auth-eof" {
socket.close(None).await.unwrap();
return;
}
let response = if outcome == "auth-denied" {
serde_json::json!({"type":"error", "code":"credential-revoked", "message":"local denial", "retryable":false})
} else {
server_adapter
.shared
.desired
.write()
.await
.as_mut()
.unwrap()
.ticket = "later-local-ticket".into();
serde_json::json!({"type":"auth.refreshed", "expiresAtMs":now_ms()+3_600_000})
};
socket
.send(Message::Text(response.to_string().into()))
.await
.unwrap();
let next = socket.next().await.unwrap().unwrap();
if outcome == "auth-denied" {
assert!(
matches!(next, Message::Close(_)),
"denied auth must not publish"
);
return;
}
let frame: serde_json::Value =
serde_json::from_str(next.to_text().unwrap()).unwrap();
assert_eq!(frame["type"], "presence.upsert");
assert_eq!(frame["device"]["ticket"], "renewed-local-ticket");
if outcome == "budget-replay" {
socket
.send(Message::Text(
serde_json::json!({
"type":"error", "code":"budget-renewal-required",
"message":"local renewal", "retryable":true,
"idempotencyKey":frame["idempotencyKey"],
})
.to_string()
.into(),
))
.await
.unwrap();
let refresh: serde_json::Value = serde_json::from_str(
socket.next().await.unwrap().unwrap().to_text().unwrap(),
)
.unwrap();
assert_eq!(refresh["type"], "auth.refresh");
socket.send(Message::Text(serde_json::json!({"type":"auth.refreshed", "expiresAtMs":now_ms()+3_600_000}).to_string().into())).await.unwrap();
let replay: serde_json::Value = serde_json::from_str(
socket.next().await.unwrap().unwrap().to_text().unwrap(),
)
.unwrap();
assert_eq!(replay, frame, "replay keeps the ticket and idempotency key");
}
assert_eq!(
server_adapter
.shared
.applied_presence
.read()
.await
.as_ref()
.unwrap()
.ticket,
"old-local-ticket"
);
let response = if outcome == "presence-denied" {
serde_json::json!({"type":"error", "code":"permission-denied", "message":"local denial", "retryable":false, "idempotencyKey":frame["idempotencyKey"]})
} else {
serde_json::json!({"type":"ack", "idempotencyKey":frame["idempotencyKey"]})
};
socket
.send(Message::Text(response.to_string().into()))
.await
.unwrap();
let next = socket.next().await.unwrap().unwrap();
if accepted {
assert!(
matches!(next, Message::Text(payload) if payload == "ping"),
"successful refresh keeps the same socket"
);
} else {
assert!(matches!(next, Message::Close(_)));
}
});
let (socket, _) = tokio_tungstenite::connect_async(endpoint.replace("http:", "ws:"))
.await
.unwrap();
let mut active = ConnectedGateway {
socket,
credential_token: "previous-local-grant".into(),
credential_expires_at_ms: now_ms() + 30_000,
lease_refresh_mode: LeaseRefreshMode::Active,
};
let result = tokio::time::timeout(
Duration::from_secs(3),
refresh_presence_authentication(&adapter, &mut active, desired.clone()),
)
.await
.unwrap();
assert_eq!(result.is_ok(), accepted, "{outcome}: {result:?}");
if accepted {
assert_eq!(*adapter.shared.applied_presence.read().await, Some(desired));
assert_eq!(active.credential_token, token);
assert!(active.credential_expires_at_ms > now_ms() + 30_000);
send_socket_keepalive(&mut active.socket).await.unwrap();
} else {
assert_eq!(*adapter.shared.applied_presence.read().await, Some(old));
if outcome != "auth-eof" {
assert!(!is_retryable_gateway_error(result.as_ref().unwrap_err()));
active.socket.close(None).await.unwrap();
}
}
tokio::time::timeout(Duration::from_secs(3), server)
.await
.unwrap()
.unwrap();
let requests = provider.requests.lock().unwrap();
assert_eq!(
requests.len(),
if outcome == "budget-replay" { 2 } else { 1 },
"only the existing bounded budget replay may request another grant"
);
if outcome == "budget-replay" {
assert_eq!(requests[1].refresh_grant, None);
assert_eq!(
requests[1].ticket_fingerprint,
requests[0].ticket_fingerprint
);
}
assert_eq!(
requests[0].refresh_grant.as_deref(),
Some("previous-local-grant")
);
assert_eq!(
requests[0].ticket_fingerprint,
base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(Sha256::digest(b"renewed-local-ticket"))
);
drop(requests);
adapter.stop().await;
}
}
#[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 authority_composition_starts_admitted_peer_owner_without_hosted_work() {
let provider = Arc::new(RecordingGrantProvider::default());
let capabilities = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
"service-device",
"native",
provider.clone(),
)
.unwrap();
let handle = capabilities
.join_room_with_architecture("authority-room", RoomArchitectureMode::Authority)
.unwrap();
let client = handle
.compose_authority_client(
crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(),
Box::new(|| None),
)
.unwrap(),
)
.await
.unwrap();
assert_eq!(
client.external_desired_peer_debug_snapshot().await["active"],
true,
"authority construction must initialize the existing Rust peer actor"
);
assert!(
client.session_registry_active(),
"admission must be closed before endpoint bind"
);
assert_eq!(
handle
.signaling
.authority_client
.lock()
.await
.as_ref()
.unwrap()
.scope,
"v2:room:authority-room",
"native and browser must gate the same capability scope"
);
assert!(
provider.requests.lock().unwrap().is_empty(),
"composition must be network-idle"
);
assert!(
handle.signaling.authority_ticket_maintenance_delay().await > Duration::from_secs(60),
"an unpublished authority must not spin a renewal timer"
);
let before = client.external_desired_peer_debug_snapshot().await;
assert!(handle
.compose_authority_client(
crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(),
Box::new(|| None)
)
.unwrap(),
)
.await
.is_err());
assert_eq!(
client.external_desired_peer_debug_snapshot().await["generation"],
before["generation"]
);
handle.close().await;
assert_eq!(
client.external_desired_peer_debug_snapshot().await["active"],
false
);
assert!(
client.session_registry_active(),
"close must not enable tokenless ingress"
);
assert!(handle
.compose_authority_client(
crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(),
Box::new(|| None)
)
.unwrap(),
)
.await
.is_err());
}
#[tokio::test]
async fn disconnected_authority_expires_admission_without_network_retry() {
authority_admission_expires_without_extra_network_work(false).await;
}
#[tokio::test]
async fn connecting_authority_expires_admission_without_restarting_grant() {
authority_admission_expires_without_extra_network_work(true).await;
}
async fn authority_admission_expires_without_extra_network_work(connecting: bool) {
let blocked = Arc::new(tokio::sync::Notify::new());
let provider = Arc::new(RecordingGrantProvider {
blocked: connecting.then(|| blocked.clone()),
..Default::default()
});
let handle = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY, "service-device", "native", provider.clone(),
).unwrap().join_room_with_architecture("expiry-room", RoomArchitectureMode::Authority).unwrap();
let client = handle.compose_authority_client(crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(), Box::new(|| None),
).unwrap()).await.unwrap();
let adapter = &handle.signaling;
let scope = "v2:room:expiry-room";
let peer = iroh::SecretKey::from_bytes(&[37; 32]).public().to_string();
let expiry = now_ms() + 150;
client.session_token_registry.update_scope_peer_admission(
scope, 1, now_ms() + 60_000, vec![peer.clone()],
).unwrap();
client.register_token_until("old-token".into(), scope.into(), 1, expiry);
client.session_token_registry.validate_and_consume_for_authenticated_peer(
"old-token", Some("expired-connection"), None, Some(&peer),
).unwrap();
client.connection_manager.upsert_pending(
"expired-connection".into(), None, Some("peer".into()), None,
).await;
client.connection_manager.set_connected("expired-connection", None).await;
let fresh_ticket = crate::session_token::expiring_ticket(
"local-endpoint-ticket", "fresh-token", scope, 1, Some(now_ms() + 900_000),
);
*adapter.shared.desired.write().await = Some(DesiredPresence {
user_id: "expiry-room".into(), local_node_id: peer,
ticket: fresh_ticket, device_name: "Service".into(), metadata: None,
ttl_ms: 60_000, online: true,
});
adapter.authority_client.lock().await.as_mut().unwrap()
.retiring_tickets.push(("old-token".into(), expiry));
let (reply, _publication) = oneshot::channel();
if connecting {
let desired = adapter.shared.desired.read().await.clone().unwrap();
adapter.commands.send(Command::Publish { desired, reply }).await.unwrap();
tokio::time::timeout(Duration::from_secs(1), async {
while provider.requests.lock().unwrap().is_empty() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
}).await.expect("actor must enter the single blocked grant request");
} else {
assert!(adapter.request_unit(|reply| Command::Patch {
patch: DevicePatch::default(), reply,
}).await.is_err());
}
let retired = tokio::time::timeout(Duration::from_secs(1), async {
loop {
let state = client.connection_manager.get_by_connection_id("expired-connection").await;
if state.is_some_and(|record| record.state == crate::connection_manager::ConnectionState::Closed) {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}).await;
let remaining = adapter.authority_client.lock().await.as_ref().unwrap().retiring_tickets.len();
let requests = provider.requests.lock().unwrap().len();
blocked.notify_one();
handle.close().await;
assert!(retired.is_ok(), "disconnected actor deferred expired admission retirement until another connect");
assert_eq!(remaining, 0);
assert_eq!(requests, usize::from(connecting), "expiry must not spend an extra network attempt");
}
#[tokio::test]
async fn authority_routes_reach_rust_once_and_withdraw_without_roster_dialing() {
let provider = Arc::new(RecordingGrantProvider::default());
let handle = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
"service-device",
"native",
provider.clone(),
)
.unwrap()
.join_room_with_architecture("authority-room", RoomArchitectureMode::Authority)
.unwrap();
let client = handle
.compose_authority_client(
crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(),
Box::new(|| None),
)
.unwrap(),
)
.await
.unwrap();
let adapter = &handle.signaling;
let now = now_ms();
let member = |id: &str| GatewayMember {
user_id: None,
device_id: id.into(),
device_name: id.into(),
platform_type: "native".into(),
metadata: None,
capabilities: None,
online: true,
updated_at_ms: now as i64,
expires_at_ms: (now + 60_000) as i64,
};
let route = |id: &str, seed: u8, scope: &str| {
let key = iroh::SecretKey::from_bytes(&[seed; 32]);
let ticket =
iroh_tickets::endpoint::EndpointTicket::new(iroh::EndpointAddr::new(key.public()))
.to_string();
GatewayDevice {
user_id: None,
device_id: id.into(),
runtime_instance_id: format!("runtime-{id}"),
node_id: key.public().to_string(),
device_name: id.into(),
platform_type: "native".into(),
ticket: crate::session_token::build_ticket(&ticket, "test-token", scope, 16),
metadata: None,
capabilities: None,
excluded_peers: Vec::new(),
online: true,
updated_at_ms: now as i64,
expires_at_ms: (now + 60_000) as i64,
}
};
handle_frame(
adapter,
ServerFrame::MembershipSnapshot {
members: vec![member("active"), member("backup"), member("roster-only")],
},
)
.await
.unwrap();
assert_eq!(
client.external_desired_peer_debug_snapshot().await["peers"],
serde_json::json!([])
);
*adapter.shared.active_grant.write().await = Some(("grant-current".into(), now + 60_000));
let lease = GatewayTopologyLease {
schema_version: 2,
topology_revision: 1,
previous_topology_revision: None,
avenue: adapter.avenue.clone(),
grant_jti: "grant-current".into(),
expires_at_ms: now + 30_000,
architecture: GatewayArchitectureLease {
policy_version: "room-architecture-v1".into(),
epoch: 1,
previous_architecture_epoch: None,
requested_mode: RoomArchitectureMode::Authority,
effective_mode: EffectiveRoomArchitecture::Authority,
phase: RoomArchitecturePhase::Settled,
reason: RoomArchitectureReason::Manual,
held_credits_microusd: 25_000,
quote_expires_at_ms: now + 30_000,
reservation_id: Some("reservation:test".into()),
},
group_encryption: None,
active: vec![route("active", 11, "v2:room:authority-room")],
admission_peers: Some(vec![
GatewayAdmissionPeer {
device_id: "active".into(),
node_id: route("active", 11, "v2:room:authority-room").node_id,
},
GatewayAdmissionPeer {
device_id: "backup".into(),
node_id: route("backup", 12, "v2:room:authority-room").node_id,
},
]),
backups: vec![route("backup", 12, "v2:room:authority-room")],
};
let mut wrong_scope = lease.clone();
let mut missing_admission = lease.clone();
missing_admission.admission_peers = None;
assert!(accept_topology_lease(adapter, missing_admission)
.await
.is_err());
wrong_scope.active = vec![route("active", 11, "user-device")];
assert!(accept_topology_lease(adapter, wrong_scope).await.is_err());
assert_eq!(*adapter.shared.topology_revision.read().await, 0);
accept_topology_lease(adapter, lease.clone()).await.unwrap();
let accepted = client.external_desired_peer_debug_snapshot().await;
assert_eq!(accepted["peers"].as_array().unwrap().len(), 1);
assert_eq!(accepted["peers"][0]["deviceId"], "active");
assert_eq!(accepted["peers"][0]["sessionId"], "runtime-active");
accept_topology_lease(adapter, lease.clone()).await.unwrap();
handle_frame(adapter, ServerFrame::Pong).await.unwrap();
assert_eq!(
client.external_desired_peer_debug_snapshot().await["revision"],
accepted["revision"]
);
handle_frame(
adapter,
ServerFrame::MembershipChanged {
operation: "delete".into(),
member: member("active"),
},
)
.await
.unwrap();
assert_eq!(
client.external_desired_peer_debug_snapshot().await["peers"],
serde_json::json!([])
);
handle_frame(
adapter,
ServerFrame::MembershipChanged {
operation: "upsert".into(),
member: member("active"),
},
)
.await
.unwrap();
assert_eq!(
client.external_desired_peer_debug_snapshot().await["peers"],
serde_json::json!([]),
"rejoining the roster must not resurrect the previous route"
);
handle.close().await;
let mut late = lease;
late.topology_revision = 2;
late.previous_topology_revision = Some(1);
assert!(
accept_topology_lease(adapter, late).await.is_err(),
"a late lease cannot restore authority admission after close"
);
assert_eq!(
client.external_desired_peer_debug_snapshot().await["active"],
false
);
assert_eq!(
client.external_desired_peer_debug_snapshot().await["peers"],
serde_json::json!([])
);
assert!(provider.requests.lock().unwrap().is_empty());
}
#[tokio::test]
async fn authority_private_routes_exchange_protected_messages_on_loopback() {
authority_loopback_payload(None, AuthorityLoopbackScenario::Renew).await;
}
#[tokio::test]
async fn authority_server_empty_routes_exchange_protected_messages_on_loopback() {
for server_is_lower in [true, false] {
authority_loopback_payload(Some(server_is_lower), AuthorityLoopbackScenario::Renew)
.await;
}
}
#[tokio::test]
async fn authority_fresh_bearer_reopens_retired_peer_without_server_dialing() {
for server_is_lower in [true, false] {
authority_loopback_payload(
Some(server_is_lower),
AuthorityLoopbackScenario::RetireAndRenew,
)
.await;
}
}
#[tokio::test]
async fn authority_gateway_revocation_denies_held_stream_and_current_route() {
for server_is_lower in [true, false] {
authority_loopback_payload(
Some(server_is_lower),
AuthorityLoopbackScenario::GatewayRevoked,
)
.await;
}
}
enum AuthorityLoopbackScenario {
Renew,
RetireAndRenew,
GatewayRevoked,
}
async fn authority_loopback_payload(
server_is_lower: Option<bool>,
scenario: AuthorityLoopbackScenario,
) {
let retire_before_renewal = matches!(scenario, AuthorityLoopbackScenario::RetireAndRenew);
let provider = Arc::new(RecordingGrantProvider::default());
let mut peers = Vec::new();
for id in ["service-a", "service-b"] {
let handle = NativeCapabilities::new(
crate::test_constants::TEST_API_KEY,
id,
"native",
provider.clone(),
)
.unwrap()
.join_room_with_architecture("loopback-room", RoomArchitectureMode::Authority)
.unwrap();
let client = handle
.compose_authority_client(
crate::Client::builder(
crate::test_constants::TEST_API_KEY.to_string(),
Box::new(|| None),
)
.unwrap()
.transport_config(crate::client::TransportConfig {
relay: false,
webrtc: None,
moq: None,
..Default::default()
}),
)
.await
.unwrap();
let endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![crate::native_node::PlutoniumProtocol::ALPN.to_vec()])
.bind_addr("127.0.0.1:0".parse::<std::net::SocketAddr>().unwrap())
.unwrap()
.bind()
.await
.unwrap();
client
.adopt_endpoint_with_router_mode(endpoint.clone(), true)
.await
.unwrap();
let ticket = client
.endpoint_ticket_with_token("v2:room:loopback-room", 16)
.await
.unwrap();
peers.push((handle, client, endpoint, ticket));
}
if let Some(server_is_lower) = server_is_lower {
peers.sort_by_key(|peer| peer.2.id());
if !server_is_lower {
peers.reverse();
}
}
let mut received_a = peers[0].1.subscribe_native_peer_data();
let mut received_b = peers[1].1.subscribe_native_peer_data();
let result = tokio::time::timeout(Duration::from_secs(15), async {
let mut leases = Vec::new();
for (local, remote) in [(0, 1), (1, 0)] {
let adapter = &peers[local].0.signaling;
let remote_adapter = &peers[remote].0.signaling;
let now = now_ms();
handle_frame(
adapter,
ServerFrame::MembershipSnapshot {
members: vec![GatewayMember {
user_id: None,
device_id: remote_adapter.device_id.clone(),
device_name: "Peer".into(),
platform_type: "native".into(),
metadata: None,
capabilities: None,
online: true,
updated_at_ms: now as i64,
expires_at_ms: (now + 60_000) as i64,
}],
},
)
.await?;
*adapter.shared.active_grant.write().await =
Some(("local-test-grant".into(), now + 60_000));
let mut lease = GatewayTopologyLease {
schema_version: 2,
topology_revision: 1,
previous_topology_revision: None,
avenue: adapter.avenue.clone(),
grant_jti: "local-test-grant".into(),
expires_at_ms: now + 30_000,
architecture: GatewayArchitectureLease {
policy_version: "room-architecture-v1".into(),
epoch: 1,
previous_architecture_epoch: None,
requested_mode: RoomArchitectureMode::Authority,
effective_mode: EffectiveRoomArchitecture::Authority,
phase: RoomArchitecturePhase::Settled,
reason: RoomArchitectureReason::Manual,
held_credits_microusd: 1,
quote_expires_at_ms: now + 30_000,
reservation_id: Some("reservation:local".into()),
},
group_encryption: None,
backups: Vec::new(),
admission_peers: Some(vec![GatewayAdmissionPeer {
device_id: remote_adapter.device_id.clone(),
node_id: peers[remote].2.id().to_string(),
}]),
active: vec![GatewayDevice {
user_id: None,
device_id: remote_adapter.device_id.clone(),
runtime_instance_id: remote_adapter.runtime_instance_id.clone(),
node_id: peers[remote].2.id().to_string(),
device_name: "Peer".into(),
platform_type: "native".into(),
ticket: peers[remote].3.clone(),
metadata: None,
capabilities: None,
excluded_peers: Vec::new(),
online: true,
updated_at_ms: now as i64,
expires_at_ms: (now + 60_000) as i64,
}],
};
if server_is_lower.is_some() && local == 0 {
lease.active.clear();
}
accept_topology_lease(adapter, lease.clone()).await?;
leases.push(lease);
}
let (a, b) = loop {
let a = peers[0]
.1
.connection_states()
.await
.into_iter()
.find(|state| state.routable);
let b = peers[1]
.1
.connection_states()
.await
.into_iter()
.find(|state| state.routable);
if let (Some(a), Some(b)) = (a, b) {
break (a, b);
}
tokio::time::sleep(Duration::from_millis(20)).await;
};
peers[0]
.1
.send_peer(&a.connection_id, b"authority-a-to-b")
.await?;
peers[1]
.1
.send_peer(&b.connection_id, b"authority-b-to-a")
.await?;
assert_eq!(received_a.recv().await?.payload, b"authority-b-to-a");
assert_eq!(received_b.recv().await?.payload, b"authority-a-to-b");
if matches!(scenario, AuthorityLoopbackScenario::GatewayRevoked) {
let (_, _, mut held, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
held.write_all(b"authorized-stream-prefix").await?;
assert!(
peers[0]
.0
.signaling
.handle_gateway_failure(&gateway_connect_error(
"local network outage",
true
),)
.await
);
held.write_all(b"still-authorized-during-outage").await?;
assert!(
peers[0]
.1
.connection_state(&a.connection_id)
.await
.unwrap()
.routable
);
let error = handle_frame(
&peers[0].0.signaling,
ServerFrame::Error {
code: "credential-revoked".into(),
message: "local administrative revocation".into(),
retryable: false,
idempotency_key: None,
metadata: Default::default(),
},
)
.await
.unwrap_err();
assert!(!peers[0].0.signaling.handle_gateway_failure(&error).await);
assert!(
held.write_all(b"revoked-held-stream").await.is_err(),
"gateway revocation must withdraw existing protected stream access"
);
assert!(peers[0]
.1
.send_peer(&a.connection_id, b"revoked-message")
.await
.is_err());
assert!(peers[0]
.1
.connection_states()
.await
.iter()
.all(|state| !state.routable));
assert!(
peers[0]
.0
.signaling
.authority_client
.lock()
.await
.as_ref()
.unwrap()
.closed
);
let mut late = leases[0].clone();
late.previous_topology_revision = Some(late.topology_revision);
late.topology_revision += 1;
assert!(
accept_topology_lease(&peers[0].0.signaling, late)
.await
.is_err(),
"late gateway projection must not reopen the revoked room owner"
);
return anyhow::Ok(());
}
let presenter = usize::from(server_is_lower.is_some());
let host = 1 - presenter;
let host_connection_id = if host == 0 {
&a.connection_id
} else {
&b.connection_id
};
let adapter = &peers[host].0.signaling;
let initial_presence = DesiredPresence {
user_id: "loopback-room".into(),
local_node_id: peers[host].2.id().to_string(),
ticket: peers[host].3.clone(),
device_name: "service-b".into(),
metadata: None,
ttl_ms: 60_000,
online: true,
};
*adapter.shared.desired.write().await = Some(initial_presence.clone());
assert!(!adapter.maintain_authority_ticket().await?);
assert!(adapter.authority_ticket_maintenance_delay().await > Duration::from_secs(60));
let (ticket, suffix) = crate::session_token::split_ticket(&peers[host].3);
let old = crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
let short_expiry = now_ms() + 30_000;
peers[host].1.register_token_until(
old.token.clone(),
old.scope.to_string(),
old.max_connections,
short_expiry,
);
adapter
.shared
.desired
.write()
.await
.as_mut()
.unwrap()
.ticket = crate::session_token::expiring_ticket(
ticket,
&old.token,
&old.scope,
old.max_connections,
Some(short_expiry),
);
assert!(adapter.authority_ticket_maintenance_delay().await <= Duration::from_millis(1));
let mut pending = adapter.shared.desired.read().await.clone().unwrap();
assert!(adapter.maintain_authority_publication(Some(&mut pending)).await?);
assert_eq!(Some(&pending), adapter.shared.desired.read().await.as_ref(),
"local renewal must preserve the pending publication acknowledgement");
let replacement = adapter
.shared
.desired
.read()
.await
.as_ref()
.unwrap()
.ticket
.clone();
assert!(
!adapter.maintain_authority_ticket().await?,
"unchanged ticket rotated twice"
);
assert_eq!(
adapter
.authority_client
.lock()
.await
.as_ref()
.unwrap()
.retiring_tickets
.len(),
1
);
let (ticket, suffix) = crate::session_token::split_ticket(&replacement);
let replacement_payload =
crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
let replacement_fingerprint =
crate::session_token::token_fingerprint(&replacement_payload.token);
if retire_before_renewal {
assert!(!peers[host].1.revoke_session_token(&old.token).is_empty());
while peers[host]
.1
.connection_manager
.get_by_connection_id(host_connection_id)
.await
.is_some()
{
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
leases[presenter].topology_revision += 1;
leases[presenter].previous_topology_revision = Some(1);
leases[presenter].active[0].ticket = replacement;
accept_topology_lease(&peers[presenter].0.signaling, leases[presenter].clone()).await?;
while peers[host]
.1
.session_token_registry
.admission_fingerprint(host_connection_id)
.as_deref()
!= Some(replacement_fingerprint.as_str())
{
tokio::time::sleep(Duration::from_millis(20)).await;
}
adapter
.authority_client
.lock()
.await
.as_mut()
.unwrap()
.retiring_tickets[0]
.1 = 0;
assert!(!adapter.maintain_authority_ticket().await?);
assert!(adapter
.authority_client
.lock()
.await
.as_ref()
.unwrap()
.retiring_tickets
.is_empty());
assert!(peers[host]
.1
.session_token_registry
.validate_and_consume(&old.token)
.is_err());
loop {
let a_ready = peers[0].1.connection_state(&a.connection_id).await
.is_some_and(|state| state.routable);
let b_ready = peers[1].1.connection_state(&b.connection_id).await
.is_some_and(|state| state.routable);
if a_ready && b_ready { break; }
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
peers[0]
.1
.connection_state(&a.connection_id)
.await
.unwrap()
.active_transport_stable_id
== a.active_transport_stable_id,
!retire_before_renewal,
"revoked physical generation must be replaced; healthy renewal must retain it"
);
assert_eq!(
peers[1]
.1
.connection_state(&b.connection_id)
.await
.unwrap()
.active_transport_stable_id
== b.active_transport_stable_id,
!retire_before_renewal
);
peers[0]
.1
.send_peer(&a.connection_id, b"after-room-token-rotation")
.await?;
assert_eq!(
received_b.recv().await?.payload,
b"after-room-token-rotation"
);
peers[1]
.1
.send_peer(&b.connection_id, b"reverse-after-room-token-rotation")
.await?;
assert_eq!(
received_a.recv().await?.payload,
b"reverse-after-room-token-rotation"
);
if server_is_lower.is_some() {
assert!(
leases[0].active.is_empty(),
"service must not dial its clients"
);
for (index, connection_id) in [(0, &a.connection_id), (1, &b.connection_id)] {
let state = peers[index]
.1
.connection_state(connection_id)
.await
.unwrap();
assert!(state.routable);
let scopes = peers[index]
.1
.connection_manager
.get_scopes(connection_id)
.await;
assert!(!scopes.iter().any(|scope| scope == "user-device"));
}
return anyhow::Ok(());
}
let unclaimed = peers[1].1.incoming_streams().await?;
assert!(
unclaimed.try_recv().is_err(),
"default messages escaped to the host"
);
let (_, _, mut send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
let mut expected = crate::stream_metadata::encode_envelope("app/custom", None)?;
expected.extend_from_slice(b"unclaimed application bytes");
send.write_all(&expected[..2]).await?;
send.write_all(&expected[2..7]).await?;
send.write_all(&expected[7..]).await?;
send.finish_and_wait_for_peer(Duration::from_secs(2))
.await?;
let incoming = unclaimed.recv().await?;
let crate::native_node::IncomingStreamType::Bi(send, recv) = incoming.stream else {
panic!("expected bidirectional application stream");
};
let (_, mut recv) = peers[1]
.1
.wrap_bi_with_prefix(
&incoming.endpoint_id,
incoming.transport_stable_id,
send,
recv,
&incoming.recv_prefix,
)
.await?;
let mut actual = Vec::new();
let mut buffer = [0; 32];
loop {
let read = recv.read(&mut buffer).await?;
if read == 0 {
break;
}
actual.extend_from_slice(&buffer[..read]);
}
assert_eq!(actual, expected);
let (_, _, mut send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
send.write_all(&crate::stream_metadata::encode_envelope(
crate::stream_metadata::DEFAULT_PEER_CHANNEL_ID,
None,
)?)
.await?;
send.write_all(&u32::MAX.to_be_bytes()).await?;
let _ = send.finish_and_wait_for_peer(Duration::from_secs(2)).await;
peers[0]
.1
.send_peer(&a.connection_id, b"after-invalid-message")
.await?;
assert_eq!(received_b.recv().await?.payload, b"after-invalid-message");
assert!(received_b.try_recv().is_err());
assert!(
unclaimed.try_recv().is_err(),
"invalid reserved channel escaped to host"
);
let protected = peers[0]
.1
.protect_outbound_application_payload(&a.connection_id, b"stale")?;
let mut frame = vec![0];
frame.extend_from_slice(&protected);
assert!(peers[1]
.1
.open_inbound_frame(
&b.connection_id,
b.remote_node_id.as_deref(),
"iroh",
Some(u64::MAX),
&frame,
)
.await
.is_err());
assert!(received_b.try_recv().is_err());
let (_, _, mut live_send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
let envelope = crate::stream_metadata::encode_envelope("app/expiry", None)?;
live_send.write_all(&envelope).await?;
let incoming = unclaimed.recv().await?;
let crate::native_node::IncomingStreamType::Bi(send, recv) = incoming.stream else {
panic!("expected bidirectional application stream");
};
let (mut expired_send, mut expired_recv) = peers[1]
.1
.wrap_bi_with_prefix(
&incoming.endpoint_id,
incoming.transport_stable_id,
send,
recv,
&incoming.recv_prefix,
)
.await?;
let mut header = vec![0; envelope.len()];
assert_eq!(expired_recv.read(&mut header).await?, envelope.len());
assert_eq!(header, envelope);
let mut denied_bytes = [0x55; 32];
{
let pending_read = expired_recv.read(&mut denied_bytes);
tokio::pin!(pending_read);
assert!(
tokio::time::timeout(Duration::from_millis(20), &mut pending_read)
.await
.is_err()
);
peers[1].1.register_token_until(
replacement_payload.token.clone(),
replacement_payload.scope.to_string(),
replacement_payload.max_connections,
0,
);
live_send.write_all(b"must-not-be-delivered").await?;
assert_eq!(
pending_read
.await
.expect_err("open stream must recheck admission after awaiting data")
.kind(),
std::io::ErrorKind::PermissionDenied
);
}
assert_eq!(
denied_bytes, [0x55; 32],
"denied reads cannot modify the consumer buffer"
);
let mut expired_recv =
expired_recv.with_plaintext_prefix(b"buffered private prefix".iter().copied());
assert_eq!(
expired_recv
.read(&mut denied_bytes)
.await
.unwrap_err()
.kind(),
std::io::ErrorKind::PermissionDenied
);
assert_eq!(
expired_send
.write_all(b"must-not-be-sent")
.await
.expect_err("open stream must recheck admission before sending")
.kind(),
std::io::ErrorKind::PermissionDenied
);
assert!(peers[1]
.1
.ensure_native_stream_admitted(&b.connection_id, b.remote_node_id.as_deref(), None,)
.is_err());
let error = peers[1]
.1
.open_inbound_frame(
&b.connection_id,
b.remote_node_id.as_deref(),
"iroh",
b.active_transport_stable_id,
&frame,
)
.await
.expect_err("expiry blocks protected messages without a gateway event");
assert!(error.to_string().contains("not admitted"), "{error:#}");
assert!(received_b.try_recv().is_err());
let recovered_ticket = peers[1]
.1
.endpoint_ticket_with_token("v2:room:loopback-room", 16)
.await?;
let (ticket, suffix) = crate::session_token::split_ticket(&recovered_ticket);
let recovered = crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
leases[0].previous_topology_revision = Some(leases[0].topology_revision);
leases[0].topology_revision += 1;
leases[0].active[0].ticket = recovered_ticket;
accept_topology_lease(&peers[0].0.signaling, leases[0].clone()).await?;
while peers[1]
.1
.session_token_registry
.admission_fingerprint(&b.connection_id)
.as_deref()
!= Some(crate::session_token::token_fingerprint(&recovered.token).as_str())
{
tokio::time::sleep(Duration::from_millis(20)).await;
}
peers[0]
.1
.send_peer(&a.connection_id, b"after-expired-admission-recovery")
.await?;
assert_eq!(
received_b.recv().await?.payload,
b"after-expired-admission-recovery"
);
assert_eq!(
expired_recv
.read(&mut denied_bytes)
.await
.unwrap_err()
.kind(),
std::io::ErrorKind::PermissionDenied,
"reauthorization cannot revive a stream that observed denial"
);
anyhow::Ok(())
})
.await;
let states_a = peers[0].1.connection_states().await;
let states_b = peers[1].1.connection_states().await;
for (handle, _, endpoint, _) in &peers {
handle.close().await;
endpoint.close().await;
}
assert!(provider.requests.lock().unwrap().is_empty());
result
.unwrap_or_else(|_| {
panic!("protected authority route timeout: server_is_lower={server_is_lower:?} a={states_a:?} b={states_b:?}")
})
.unwrap();
}
#[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 authority_return_retry_never_resets_budget_or_retries_terminal_errors() {
let missing =
gateway_server_error("room-authority-unavailable".into(), "local".into(), true);
let mut available = false;
assert_eq!(
gateway_reconnect_delay(10, &missing, &mut available),
Some(MAX_RECONNECT_DELAY),
"initial publication cannot accelerate a host that has never been available"
);
available = true;
for attempt in 1..=13 {
let error = if attempt < 10 {
gateway_connect_error("network unavailable", true)
} else {
gateway_server_error("room-authority-unavailable".into(), "local".into(), true)
};
let expected = match attempt {
10 => Some(reconnect_delay(1)),
12.. => None,
_ => Some(reconnect_delay(attempt)),
};
assert_eq!(
gateway_reconnect_delay(attempt, &error, &mut available),
expected
);
}
assert!(!available);
for (code, retryable) in [
("credential-revoked", false),
("room-architecture-conflict", false),
("room-authority-unavailable", false),
("budget-renewal-required", true),
] {
available = true;
let error = gateway_server_error(code.into(), "local".into(), retryable);
assert_eq!(
gateway_reconnect_delay(10, &error, &mut available),
retryable.then_some(MAX_RECONNECT_DELAY)
);
assert!(
available,
"unrelated/permanent failures must not consume the return allowance"
);
}
available = true;
assert_eq!(gateway_reconnect_delay(12, &missing, &mut available), None);
assert!(
available,
"even an unused allowance cannot bypass the attempt cap"
);
let untyped = gateway_connect_error("room-authority-unavailable", true);
assert_eq!(
gateway_reconnect_delay(10, &untyped, &mut available),
Some(MAX_RECONNECT_DELAY)
);
}
#[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
#[tokio::test]
async fn recovered_gateway_retries_host_return_without_another_capped_delay() {
gateway_close_actor_scenario(false).await;
}
#[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
#[tokio::test]
async fn revoked_gateway_socket_stops_authority_without_minting_another_grant() {
gateway_close_actor_scenario(true).await;
}
#[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
async fn gateway_close_actor_scenario(revoked: bool) {
struct OutageProvider {
inner: RecordingGrantProvider,
failures: std::sync::atomic::AtomicU8,
calls: std::sync::atomic::AtomicU8,
}
#[async_trait]
impl NativeGatewayGrantProvider for OutageProvider {
async fn grant(
&self,
request: NativeGatewayGrantRequest,
) -> Result<NativeGatewayGrant> {
self.calls.fetch_add(1, Ordering::SeqCst);
if self
.failures
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
remaining.checked_sub(1)
})
.is_ok()
{
bail!("local injected control-plane outage");
}
self.inner.grant(request).await
}
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let endpoint = format!("http://{}", listener.local_addr().unwrap());
let provider = Arc::new(OutageProvider {
inner: RecordingGrantProvider {
gateway_url: Some(endpoint.clone()),
token: Some(format!(
"local.{}.test",
base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(br#"{"jti":"local-outage-grant"}"#)
)),
..Default::default()
},
failures: 0.into(),
calls: 0.into(),
});
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint,
app_tag: "app_native_test".into(),
device_id: "service-consumer".into(),
platform_type: "native".into(),
avenue: NativeCoordinationAvenue {
kind: "room".into(),
id: "local-room".into(),
},
architecture: Some(RoomArchitectureMode::Authority),
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,
})
.unwrap();
let (cut, cut_rx) = oneshot::channel();
let (host_missing, host_missing_rx) = oneshot::channel();
let server = tokio::spawn(async move {
let mut cut_rx = Some(cut_rx);
let mut host_missing = Some(host_missing);
for connection in 0..3 {
let (tcp, _) = listener.accept().await.unwrap();
let mut socket = tokio_tungstenite::accept_hdr_async(tcp,
|_: &tokio_tungstenite::tungstenite::handshake::server::Request,
mut response: tokio_tungstenite::tungstenite::handshake::server::Response| {
response.headers_mut().insert("Sec-WebSocket-Protocol",
HeaderValue::from_static(GATEWAY_PROTOCOL));
Ok(response)
}).await.unwrap();
let auth: serde_json::Value =
serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
.unwrap();
assert_eq!(auth["type"], "auth");
let response = if connection == 1 {
serde_json::json!({"type":"error", "code":"room-authority-unavailable",
"message":"Host is returning", "retryable":true})
} else {
serde_json::json!({"type":"ready", "budgetRemainingMicrousd":1000})
};
socket
.send(Message::Text(response.to_string().into()))
.await
.unwrap();
if connection == 0 {
cut_rx.take().unwrap().await.unwrap();
let frame =
revoked.then(|| tokio_tungstenite::tungstenite::protocol::CloseFrame {
code: 4403.into(),
reason: "credential revoked".into(),
});
socket.close(frame).await.unwrap();
if revoked {
return;
}
} else if connection == 1 {
host_missing.take().unwrap().send(()).unwrap();
socket.close(None).await.unwrap();
} else {
while let Some(Ok(frame)) = socket.next().await {
match frame {
Message::Text(text) if text == "ping" => {
socket.send(Message::Text("pong".into())).await.unwrap();
}
Message::Ping(payload) => {
socket.send(Message::Pong(payload)).await.unwrap();
}
Message::Close(_) => break,
_ => {}
}
}
}
}
});
tokio::time::timeout(
Duration::from_secs(3),
adapter.update_presence(
"local-room",
&iroh::SecretKey::from_bytes(&[17; 32]).public().to_string(),
"local-ticket",
true,
"Consumer",
900_000,
None,
),
)
.await
.unwrap()
.unwrap();
if revoked {
*adapter.authority_client.lock().await = Some(AuthorityPeerBinding {
client: std::sync::Weak::new(),
scope: "v2:room:local-room".into(),
revision: 1,
admission_revision: 1,
last_payload: "[]".into(),
retiring_tickets: Vec::new(),
closed: false,
});
cut.send(()).unwrap();
let stopped = tokio::time::timeout(Duration::from_secs(3), async {
while !adapter
.authority_client
.lock()
.await
.as_ref()
.unwrap()
.closed
{
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await;
adapter.stop().await;
server.abort();
let _ = server.await;
assert!(
stopped.is_ok(),
"revocation close did not reach authority owner"
);
assert_eq!(
provider.calls.load(Ordering::SeqCst),
1,
"revoked socket must not mint a replacement grant"
);
return;
}
provider.failures.store(9, Ordering::SeqCst);
cut.send(()).unwrap();
tokio::time::timeout(Duration::from_secs(3), async {
while provider.calls.load(Ordering::SeqCst) < 2 {
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await
.unwrap();
tokio::time::pause();
for attempt in 2..=9 {
tokio::time::advance(reconnect_delay(attempt - 1) + Duration::from_millis(2)).await;
for _ in 0..100 {
tokio::task::yield_now().await;
}
assert_eq!(provider.calls.load(Ordering::SeqCst), attempt + 1);
}
tokio::time::advance(reconnect_delay(9) - Duration::from_millis(1)).await;
tokio::time::resume();
tokio::time::timeout(Duration::from_secs(3), host_missing_rx)
.await
.unwrap()
.unwrap();
let recovered = tokio::time::timeout(Duration::from_secs(3), async {
while adapter.shared.applied_presence.read().await.is_none() {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await;
adapter.stop().await;
server.abort();
let _ = server.await;
assert!(
recovered.is_ok(),
"host return must not inherit another 30-second network backoff"
);
assert_eq!(
provider.calls.load(Ordering::SeqCst),
12,
"initial connection plus eleven recovery attempts; never reset the outage budget"
);
}
#[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_gateway_revocation_is_terminal_authority_withdrawal() {
let revoked = decode_frame(Message::Close(Some(
tokio_tungstenite::tungstenite::protocol::CloseFrame {
code: 4403.into(),
reason: "credential revoked".into(),
},
)))
.unwrap_err();
assert!(is_gateway_revocation(&revoked));
assert!(!is_retryable_gateway_error(&revoked));
assert!(!is_gateway_revocation(
&decode_frame(Message::Close(None)).unwrap_err()
));
assert!(!is_gateway_revocation(&gateway_connect_error(
"credential-revoked",
false
)));
assert!(!is_gateway_revocation(&gateway_server_error(
"budget-exhausted".into(),
"local".into(),
false
)));
let wire = gateway_server_error("credential-revoked".into(), "local".into(), true);
assert!(
!is_retryable_gateway_error(&wire),
"revocation cannot grant a retry through a contradictory flag"
);
assert!(is_gateway_revocation(&wire.context("operation failed")));
let certificate = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
status: 412,
message: "device revoked".into(),
reason: Some("device-certificate-revoked".into()),
})
.context("obtain native coordination grant");
assert!(is_gateway_revocation(&certificate));
}
#[tokio::test]
async fn command_reply_preserves_revocation_for_the_socket_owner() {
let (reply, result) = oneshot::channel::<Result<()>>();
let error = complete_command(
reply,
Err(gateway_server_error(
"credential-revoked".into(),
"local".into(),
false,
)),
"session.put",
)
.unwrap_err();
let caller = result.await.unwrap().unwrap_err();
assert!(is_gateway_revocation(&caller));
assert!(!is_retryable_gateway_error(&caller));
assert!(is_gateway_revocation(&error));
assert!(!is_retryable_gateway_error(&error));
}
#[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,
));
}
#[tokio::test]
async fn command_failure_reply_preserves_the_native_error_chain() {
let (reply, response) = oneshot::channel();
let returned = complete_command::<()>(
reply,
Err(anyhow!("control-plane detail").context("renewal failed")),
"native device patch",
)
.expect_err("the actor must observe the operation failure");
let caller = response
.await
.expect("reply")
.expect_err("the caller must observe the operation failure");
assert_eq!(caller.to_string(), "renewal failed: control-plane detail");
assert!(
format!("{returned:#}").contains("renewal failed: control-plane detail"),
"the actor should retain the same detailed chain"
);
}
#[tokio::test]
async fn command_failure_reply_preserves_service_metadata_for_both_owners() {
for code in crate::service_errors::SERVICE_ERROR_CODES {
let metadata = serde_json::json!({
"code": code, "retryable": false, "scope": "app",
"operation": "presence.upsert", "requestId": "request-42",
"retryAfterMs": 120000, "resetAt": 1800000000000_u64,
});
for http in [false, true] {
let denial = if http {
anyhow::Error::new(ControlPlaneHttpError {
status: 429, message: "denied".into(), reason: None,
service_error: crate::service_errors::ServiceError::from_value(&metadata),
})
} else {
gateway_server_error_metadata(code.to_string(), "denied".into(), false,
metadata.as_object().unwrap().clone())
}.context("refresh failed");
let (reply, response) = oneshot::channel();
let actor = complete_command::<()>(reply, Err(denial), "presence.upsert").unwrap_err();
let caller = response.await.unwrap().unwrap_err();
for error in [&actor, &caller] {
let service = crate::service_errors::ServiceError::from_error(error).unwrap();
assert_eq!(service.code, *code);
assert_eq!(service.request_id.as_deref(), Some("request-42"));
assert_eq!(service.retry_after_ms, Some(120000));
assert_eq!(service.reset_at, Some(1800000000000));
assert!(!is_retryable_gateway_error(error));
}
}
}
}
#[tokio::test]
async fn terminal_publication_preserves_service_denial_without_another_grant() {
struct DeniedProvider(std::sync::atomic::AtomicUsize);
#[async_trait]
impl NativeGatewayGrantProvider for DeniedProvider {
async fn grant(&self, _: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Err(anyhow::Error::new(ControlPlaneHttpError {
status: 429, message: "app cap reached".into(), reason: None,
service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-budget-exhausted", "retryable": false, "scope": "app",
"operation": "gateway.grant.issue", "requestId": "terminal-request",
})),
}))
}
}
let provider = Arc::new(DeniedProvider(0.into()));
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
device_id: "device-1".into(), platform_type: "desktop".into(),
avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
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,
}).unwrap();
let desired = DesiredPresence {
user_id: "principal-1".into(), local_node_id: "node-1".into(),
ticket: "ticket-1".into(), device_name: "Native".into(),
metadata: None, ttl_ms: 900000, online: true,
};
for _ in 0..3 {
let (reply, response) = oneshot::channel();
adapter.commands.send(Command::Publish { desired: desired.clone(), reply }).await.unwrap();
let error = tokio::time::timeout(Duration::from_secs(2), response).await.unwrap().unwrap().unwrap_err();
let service = crate::service_errors::ServiceError::from_error(&error).unwrap();
assert_eq!(service.code, "app-budget-exhausted");
assert_eq!(service.request_id.as_deref(), Some("terminal-request"));
assert!(!is_retryable_gateway_error(&error));
}
assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
adapter.stop().await;
}
#[test]
fn gateway_wire_service_errors_preserve_safe_metadata_and_cooldown() {
for code in crate::service_errors::SERVICE_ERROR_CODES {
let frame = decode_frame(Message::Text(serde_json::json!({
"type": "error", "code": code, "message": "admission denied",
"retryable": true, "retryAfterMs": 120000,
"scope": "app", "operation": "gateway.grant.issue",
"requestId": "request-1", "providerCost": 99,
"documentationUrl": "https://untrusted.invalid"
}).to_string().into())).unwrap().unwrap();
let ServerFrame::Error { code: actual, message, retryable, metadata, .. } = frame else {
panic!("expected service denial");
};
let error = gateway_server_error_metadata(actual, message, retryable, metadata)
.context("native gateway admission");
let service = gateway_service_error(&error).unwrap();
assert_eq!(service.code, *code);
assert_eq!(service.scope.as_deref(), Some("app"));
assert_eq!(service.request_id.as_deref(), Some("request-1"));
let serialized = serde_json::to_value(service).unwrap();
assert!(serialized.get("providerCost").is_none());
assert_eq!(serialized["documentationUrl"], crate::service_errors::ERROR_DOCUMENTATION_URL);
assert_eq!(gateway_reconnect_delay(1, &error, &mut true), Some(Duration::from_secs(120)));
}
}
#[test]
fn legacy_edge_frames_preserve_canonical_scope_and_owner_cooldown() {
for (legacy, scope) in [("edge-rate-exceeded", "ip"), ("route-rate-exceeded", "avenue")] {
let frame = decode_frame(Message::Text(serde_json::json!({
"type": "error", "code": legacy, "message": "rate limited", "retryable": true,
"retryAfterMs": 120000, "requestId": "legacy-1"
}).to_string().into())).unwrap().unwrap();
let ServerFrame::Error { code, message, retryable, metadata, .. } = frame else {
panic!("expected legacy denial");
};
let error = gateway_server_error_metadata(code, message, retryable, metadata);
let service = gateway_service_error(&error).unwrap();
assert_eq!(service.code, "edge-rate-limited");
assert_eq!(service.scope.as_deref(), Some(scope));
assert_eq!(service.request_id.as_deref(), Some("legacy-1"));
assert_eq!(gateway_reconnect_delay(1, &error, &mut true), Some(Duration::from_secs(120)));
}
}
#[tokio::test]
async fn service_error_observation_is_capability_local_and_preserves_retryability() {
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
device_id: "device-1".into(), platform_type: "desktop".into(),
avenue: NativeCoordinationAvenue { kind: "space".into(), id: "space-1".into() },
architecture: None, room_delivery: None,
grant_provider: Arc::new(RecordingGrantProvider::default()),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
}).unwrap();
let mut observations = adapter.subscribe_service_errors();
let isolated = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
device_id: "device-2".into(), platform_type: "desktop".into(),
avenue: NativeCoordinationAvenue { kind: "space".into(), id: "space-2".into() },
architecture: None, room_delivery: None,
grant_provider: Arc::new(RecordingGrantProvider::default()),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
}).unwrap();
let mut isolated_observations = isolated.subscribe_service_errors();
let known = anyhow::Error::new(ControlPlaneHttpError {
status: 429, message: "limited".into(), reason: None,
service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-rate-limited", "retryable": true,
"scope": "app", "operation": "gateway.grant.issue",
"requestId": "request-1", "providerCost": 99,
})),
}).context("wrapped gateway denial");
assert!(adapter.handle_gateway_failure(&known).await);
let observation = observations.recv().await.unwrap();
assert_eq!(observation.avenue.kind, "space");
assert_eq!(observation.avenue.id, "space-1");
assert_eq!(observation.runtime_instance_id, adapter.runtime_instance_id);
assert_eq!(observation.service_error.code, "app-rate-limited");
assert_eq!(observation.service_error.request_id.as_deref(), Some("request-1"));
let serialized = serde_json::to_value(&observation).unwrap();
assert!(serialized.get("avenue").is_some());
assert!(serialized.get("runtimeInstanceId").is_some());
assert!(serialized.get("serviceError").is_some());
assert!(serialized.get("runtime_instance_id").is_none());
assert!(serialized["serviceError"].get("providerCost").is_none());
assert!(serialized["serviceError"].get("message").is_none());
assert!(matches!(isolated_observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
let terminal = anyhow::Error::new(ControlPlaneHttpError {
status: 403, message: "revoked".into(), reason: None,
service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-budget-exhausted", "retryable": false,
"scope": "avenue", "operation": "gateway.grant.issue",
})),
}).context("wrapped terminal denial");
assert!(!adapter.handle_gateway_failure(&terminal).await);
let terminal_observation = observations.recv().await.unwrap();
assert!(!terminal_observation.service_error.retryable);
assert!(matches!(isolated_observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
let unknown = anyhow!("connection reset");
assert!(adapter.handle_gateway_failure(&unknown).await);
assert!(matches!(observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
adapter.stop().await;
isolated.stop().await;
}
#[tokio::test]
async fn closed_capability_rejects_service_error_subscription() {
let handle = NativeCapabilityHandle::new(GatewayOptions {
endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
device_id: "device-1".into(), platform_type: "desktop".into(),
avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
architecture: None, room_delivery: None,
grant_provider: Arc::new(RecordingGrantProvider::default()),
#[cfg(feature = "managed-group-encryption")]
managed_group_signer: None,
#[cfg(feature = "managed-group-encryption")]
managed_group_store: None,
}).unwrap();
handle.close().await;
assert!(handle.subscribe_service_errors().is_err());
}
#[tokio::test]
async fn native_actor_republication_cannot_bypass_service_cooldown() {
struct LimitedProvider(std::sync::atomic::AtomicUsize);
#[async_trait]
impl NativeGatewayGrantProvider for LimitedProvider {
async fn grant(&self, _: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Err(anyhow::Error::new(ControlPlaneHttpError {
status: 429, message: "try later".into(), reason: None,
service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
"code": "app-rate-limited", "retryable": true, "retryAfterMs": 120000
})),
}))
}
}
let provider = Arc::new(LimitedProvider(0.into()));
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
device_id: "device-1".into(), platform_type: "desktop".into(),
avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
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,
}).unwrap();
tokio::time::pause();
let desired = DesiredPresence {
user_id: "principal-1".into(), local_node_id: "node-1".into(),
ticket: "ticket-1".into(), device_name: "Native".into(),
metadata: None, ttl_ms: 900000, online: true,
};
let (reply, first) = oneshot::channel();
adapter.commands.send(Command::Publish { desired: desired.clone(), reply }).await.unwrap();
for _ in 0..20 { tokio::task::yield_now().await; }
tokio::time::advance(Duration::from_millis(1)).await;
for _ in 0..20 { tokio::task::yield_now().await; }
assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
tokio::time::advance(Duration::from_secs(60)).await;
let (reply, second) = oneshot::channel();
adapter.commands.send(Command::Publish { desired, reply }).await.unwrap();
assert!(first.await.unwrap().is_err(), "old publication is superseded");
for _ in 0..20 { tokio::task::yield_now().await; }
assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
tokio::time::advance(Duration::from_secs(59)).await;
for _ in 0..20 { tokio::task::yield_now().await; }
assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
tokio::time::advance(Duration::from_millis(1001)).await;
for _ in 0..20 { tokio::task::yield_now().await; }
assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 2);
adapter.stop().await;
assert!(second.await.unwrap().is_err());
}
#[test]
fn service_cooldown_is_a_minimum_not_capped_network_backoff() {
let denial = |metadata: serde_json::Value| {
anyhow::Error::new(ControlPlaneHttpError {
status: 429,
message: "try later".into(),
reason: None,
service_error: crate::service_errors::ServiceError::from_value(&metadata),
})
.context("obtain native coordination grant")
};
let error = denial(serde_json::json!({
"code": "app-rate-limited", "retryable": true, "retryAfterMs": 120000
}));
let mut allowance = true;
assert_eq!(gateway_service_retry_delay(&error), Duration::from_secs(120));
assert_eq!(gateway_reconnect_delay(1, &error, &mut allowance), Some(Duration::from_secs(120)));
assert_eq!(gateway_reconnect_delay(10, &error, &mut allowance), Some(Duration::from_secs(120)));
assert_eq!(gateway_reconnect_delay(MAX_RECONNECT_ATTEMPTS, &error, &mut allowance), None);
assert!(allowance, "financial denial cannot consume the host return allowance");
let reset = now_ms() + 300000;
let error = denial(serde_json::json!({
"code": "principal-rate-limited", "retryable": true,
"retryAfterMs": 120000, "resetAt": reset
}));
let before = now_ms();
let delay = gateway_service_retry_delay(&error);
assert!(delay >= Duration::from_millis(reset.saturating_sub(now_ms())));
assert!(delay <= Duration::from_millis(reset.saturating_sub(before)));
let terminal = denial(serde_json::json!({
"code": "credit-exhausted", "retryable": false, "retryAfterMs": 120000
}));
assert_eq!(gateway_reconnect_delay(1, &terminal, &mut allowance), None);
let legacy = gateway_connect_error("network outage", true);
assert_eq!(gateway_service_retry_delay(&legacy), Duration::ZERO);
assert_eq!(gateway_reconnect_delay(2, &legacy, &mut allowance), Some(reconnect_delay(2)));
}
#[test]
fn terminal_control_plane_rejections_are_not_retried() {
let revoked = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
status: 403,
message: "device revoked".into(),
reason: Some("device-certificate-revoked".into()),
})
.context("obtain native coordination grant");
let throttled = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
status: 429,
message: "try later".into(),
reason: None,
})
.context("obtain native coordination grant");
let unavailable = anyhow::Error::new(ControlPlaneHttpError {
service_error: None,
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() {
let (mut sender, mut receiver) = futures::channel::mpsc::channel(1);
send_socket_keepalive(&mut sender).await.unwrap();
assert!(
matches!(receiver.next().await, Some(Message::Text(payload)) if payload == "ping")
);
assert!(decode_frame(Message::Text("pong".into())).unwrap().is_none());
}
#[cfg(feature = "testing-endpoints")]
#[tokio::test]
async fn busy_socket_keepalive() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let endpoint = format!("http://{}", listener.local_addr().unwrap());
let provider = Arc::new(RecordingGrantProvider {
gateway_url: Some(endpoint.clone()),
token: Some(format!("local.{}.test",
base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(br#"{"jti":"busy-socket-grant"}"#))),
..Default::default()
});
let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
endpoint,
app_tag: "app_native_test".into(),
device_id: "busy-device".into(),
platform_type: "native".into(),
avenue: NativeCoordinationAvenue { kind: "user".into(), id: "busy-user".into() },
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,
}).unwrap();
let server = tokio::spawn(async move {
let (tcp, _) = listener.accept().await.unwrap();
let mut socket = tokio_tungstenite::accept_hdr_async(tcp,
|_: &tokio_tungstenite::tungstenite::handshake::server::Request,
mut response: tokio_tungstenite::tungstenite::handshake::server::Response| {
response.headers_mut().insert("Sec-WebSocket-Protocol",
HeaderValue::from_static(GATEWAY_PROTOCOL));
Ok(response)
}).await.unwrap();
let auth: serde_json::Value = serde_json::from_str(
socket.next().await.unwrap().unwrap().to_text().unwrap()).unwrap();
assert_eq!(auth["socketLiveness"], "ping-v1");
socket.send(Message::Text(
serde_json::json!({"type":"ready", "budgetRemainingMicrousd":1000})
.to_string().into())).await.unwrap();
let start = tokio::time::Instant::now();
let mut traffic = tokio::time::interval(Duration::from_millis(100));
let mut pings = 0;
loop {
tokio::select! {
_ = traffic.tick() => {
socket.send(Message::Text("pong".into())).await.unwrap();
}
frame = socket.next() => {
assert!(matches!(frame, Some(Ok(Message::Text(payload))) if payload == "ping"),
"keepalive must not republish presence or refresh authorization");
pings += 1;
if pings >= 2 && start.elapsed() >= SOCKET_KEEPALIVE_INTERVAL / 2 {
return;
}
}
}
}
});
let published = tokio::time::timeout(Duration::from_secs(3), adapter.update_presence(
"busy-user", &iroh::SecretKey::from_bytes(&[19; 32]).public().to_string(),
"busy-ticket", true, "Busy", 900_000, None)).await;
let mut server = server;
let observed = tokio::time::timeout(SOCKET_KEEPALIVE_INTERVAL * 2 + Duration::from_secs(5),
&mut server).await;
adapter.stop().await;
if observed.is_err() {
server.abort();
let _ = server.await;
}
published.expect("initial presence timed out").unwrap();
observed.expect("incoming traffic starved native keepalive").unwrap();
assert_eq!(provider.requests.lock().unwrap().len(), 1,
"keepalive must reuse the authenticated socket and grant");
}
#[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,
admission_peers: None,
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,
admission_peers: None,
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,
admission_peers: None,
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());
}
}