Skip to main content

openrtc/
native_coordination_gateway.rs

1//! Provider-neutral native coordination-gateway adapter.
2//!
3//! The adapter owns only authenticated gateway delivery: credential exchange,
4//! WebSocket subscription, bounded reconnect, presence publication, and
5//! signaling event projection. Rust `Client` remains the sole peer lifecycle
6//! and transport authority.
7
8#[cfg(feature = "managed-group-encryption")]
9use crate::managed_group_controller::{
10    ManagedArtifactChunk, ManagedGroupAction, ManagedGroupController, ManagedPreparePage,
11};
12#[cfg(feature = "managed-group-encryption")]
13use crate::native::DeviceSigner;
14use crate::native::{
15    ControlPlaneHttpError, EffectiveRoomArchitecture, RoomArchitectureMode, RoomArchitecturePhase,
16    RoomArchitectureReason, RoomArchitectureSnapshot, RoomDelivery,
17};
18use crate::signaling::{
19    Device, DeviceCapabilities, DeviceEvent, SessionEvent, SignalingBackend, SignalingEnvelope,
20    SignalingSession,
21};
22use anyhow::{anyhow, bail, Context, Result};
23use async_trait::async_trait;
24use base64::Engine as _;
25use futures::{stream::BoxStream, Sink, SinkExt, StreamExt};
26use serde::{Deserialize, Serialize};
27use sha2::{Digest, Sha256};
28use std::collections::HashMap;
29use std::error::Error as StdError;
30use std::fmt;
31use std::str::FromStr;
32use std::sync::atomic::{AtomicBool, Ordering};
33use std::sync::{Arc, Mutex};
34use std::time::{Duration, SystemTime, UNIX_EPOCH};
35use tokio::sync::{broadcast, mpsc, oneshot, watch, RwLock};
36use tokio_tungstenite::tungstenite::{
37    client::IntoClientRequest,
38    http::{HeaderValue, Uri},
39    Message,
40};
41use uuid::Uuid;
42#[cfg(feature = "managed-group-encryption")]
43use zeroize::Zeroize;
44
45const WIRE_PROTOCOL_VERSION: u8 = 1;
46const GRANT_PROTOCOL_VERSION: u8 = 2;
47const GATEWAY_PROTOCOL: &str = "openrtc.v2";
48const GATEWAY_AUTH_PROTOCOL_PREFIX: &str = "openrtc.auth.";
49const REQUEST_TIMEOUT: Duration = Duration::from_secs(15);
50const ACK_TIMEOUT: Duration = Duration::from_secs(15);
51const AUTH_REFRESH_SKEW_MS: u64 = 5 * 60_000;
52const AUTHORITY_TICKET_RENEWAL_SKEW_MS: u64 = 60_000;
53// Keep an otherwise idle native connection alive below common 30-60 minute
54// NAT/proxy idle ceilings. Cloudflare handles protocol ping frames at the
55// WebSocket edge without waking a hibernating Durable Object, so this does not
56// create a presence write, logical usage event, or recurring DO execution.
57const SOCKET_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30);
58#[cfg(feature = "adaptive-room-sentinel")]
59pub(crate) const MAX_RECONNECT_ATTEMPTS: u8 = 12;
60#[cfg(not(feature = "adaptive-room-sentinel"))]
61const MAX_RECONNECT_ATTEMPTS: u8 = 12;
62#[cfg(feature = "adaptive-room-sentinel")]
63pub(crate) const MAX_OPERATION_ATTEMPTS: usize = 2;
64#[cfg(not(feature = "adaptive-room-sentinel"))]
65const MAX_OPERATION_ATTEMPTS: usize = 2;
66const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30);
67const MAX_EXCLUDED_PEERS: usize = 250;
68const MAX_AUTHORITY_PRIORITY_PEERS: usize = 8;
69const MAX_AUTHORITY_RELEVANT_ENTITIES: usize = 32;
70#[cfg(feature = "managed-group-encryption")]
71const MAX_MANAGED_BATCH_MESSAGES: usize = 100;
72pub const OPENRTC_PRODUCTION_COORDINATION_GATEWAY: &str = "https://gateway.openrtc.app";
73
74#[derive(Debug, Clone)]
75pub struct GatewayConnectError {
76    message: String,
77    retryable: bool,
78    code: Option<String>,
79    pub service_error: Option<crate::service_errors::ServiceError>,
80}
81
82impl fmt::Display for GatewayConnectError {
83    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
84        formatter.write_str(&self.message)
85    }
86}
87
88impl StdError for GatewayConnectError {}
89
90// An operation failure is delivered both to its caller and the socket owner,
91// and a parked publication can be requested again. Keep the typed denial in
92// each copy rather than reducing it to a message at these fan-out boundaries.
93fn copy_gateway_failure(error: &anyhow::Error) -> anyhow::Error {
94    let detail = format!("{error:#}");
95    if let Some(denial) = error.downcast_ref::<GatewayConnectError>() {
96        return anyhow::Error::new(denial.clone()).context(detail);
97    }
98    if let Some(denial) = error.downcast_ref::<ControlPlaneHttpError>() {
99        return anyhow::Error::new(ControlPlaneHttpError {
100            status: denial.status,
101            message: denial.message.clone(),
102            reason: denial.reason.clone(),
103            service_error: denial.service_error.clone(),
104        }).context(detail);
105    }
106    anyhow!(detail)
107}
108
109fn gateway_connect_error(message: impl Into<String>, retryable: bool) -> anyhow::Error {
110    anyhow!(GatewayConnectError {
111        message: message.into(),
112        retryable,
113        code: None,
114        service_error: None,
115    })
116}
117
118fn gateway_server_error(code: String, message: String, retryable: bool) -> anyhow::Error {
119    gateway_server_error_metadata(code, message, retryable, Default::default())
120}
121
122fn gateway_server_error_metadata(
123    code: String,
124    message: String,
125    retryable: bool,
126    mut metadata: serde_json::Map<String, serde_json::Value>,
127) -> anyhow::Error {
128    metadata.insert("code".into(), serde_json::Value::String(code.clone()));
129    metadata.insert("retryable".into(), serde_json::Value::Bool(retryable));
130    let service_error = crate::service_errors::ServiceError::from_value(&metadata.into());
131    anyhow!(GatewayConnectError {
132        message: format!("coordination gateway {code}: {message}"),
133        retryable,
134        code: Some(code),
135        service_error,
136    })
137}
138
139fn is_retryable_gateway_error(error: &anyhow::Error) -> bool {
140    if is_gateway_revocation(error) {
141        return false;
142    }
143    if let Some(classified) = error.downcast_ref::<GatewayConnectError>() {
144        return classified.retryable;
145    }
146    if let Some(control_plane) = error
147        .chain()
148        .find_map(|source| source.downcast_ref::<ControlPlaneHttpError>())
149    {
150        return control_plane.is_retryable();
151    }
152    true
153}
154
155fn gateway_service_error(error: &anyhow::Error) -> Option<&crate::service_errors::ServiceError> {
156    crate::service_errors::ServiceError::from_error(error)
157}
158
159fn gateway_service_retry_delay(error: &anyhow::Error) -> Duration {
160    let Some(service) = gateway_service_error(error) else {
161        return Duration::ZERO;
162    };
163    Duration::from_millis(
164        service.retry_after_ms.unwrap_or(0).max(
165            service.reset_at.unwrap_or(0).saturating_sub(now_ms()),
166        ),
167    )
168}
169
170fn is_gateway_revocation(error: &anyhow::Error) -> bool {
171    error
172        .downcast_ref::<GatewayConnectError>()
173        .is_some_and(|error| error.code.as_deref() == Some("credential-revoked"))
174        || error
175            .downcast_ref::<ControlPlaneHttpError>()
176            .is_some_and(|error| error.reason.as_deref() == Some("device-certificate-revoked"))
177}
178
179fn operation_requires_budget_renewal(code: &str, retryable: bool) -> bool {
180    retryable && code == "budget-renewal-required"
181}
182
183fn terminal_refresh_error(
184    lease_error: anyhow::Error,
185    fallback_error: Option<anyhow::Error>,
186) -> anyhow::Error {
187    fallback_error.unwrap_or(lease_error)
188}
189
190#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
191#[serde(rename_all = "camelCase")]
192pub struct NativeCoordinationAvenue {
193    pub kind: String,
194    pub id: String,
195}
196
197/// A historical, capability-local service denial observation. This is a
198/// diagnostic projection only; it does not authorize recovery or change
199/// gateway retry policy.
200#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
201#[serde(rename_all = "camelCase")]
202pub struct NativeServiceErrorObservation {
203    pub avenue: NativeCoordinationAvenue,
204    pub runtime_instance_id: String,
205    pub service_error: crate::service_errors::ServiceError,
206}
207
208#[derive(Debug, Clone, PartialEq, Eq)]
209pub struct NativeGatewayGrantRequest {
210    pub avenue: NativeCoordinationAvenue,
211    pub device_id: String,
212    pub runtime_instance_id: String,
213    pub ticket_fingerprint: String,
214    pub purpose: String,
215    pub architecture: Option<RoomArchitectureMode>,
216    pub room_delivery: Option<RoomDelivery>,
217    pub refresh_grant: Option<String>,
218}
219
220#[async_trait]
221pub trait NativeGatewayGrantProvider: Send + Sync {
222    async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant>;
223}
224
225#[derive(Debug, Clone, Deserialize)]
226#[serde(rename_all = "camelCase")]
227pub struct NativeGatewayGrant {
228    pub protocol_version: u8,
229    pub gateway_url: String,
230    pub route_key: String,
231    pub token: String,
232    pub expires_at_ms: u64,
233}
234
235/// Configuration for a pure-native OpenRTC 2.0 consumer. Identity assertion,
236/// device proof, attestation, and secure storage remain host responsibilities;
237/// the provider returns only OpenRTC-issued, avenue-bound grants.
238pub struct GatewayOptions {
239    pub endpoint: String,
240    pub app_tag: String,
241    pub device_id: String,
242    pub platform_type: String,
243    pub avenue: NativeCoordinationAvenue,
244    pub architecture: Option<RoomArchitectureMode>,
245    pub room_delivery: Option<RoomDelivery>,
246    pub grant_provider: Arc<dyn NativeGatewayGrantProvider>,
247    #[cfg(feature = "managed-group-encryption")]
248    pub managed_group_signer: Option<Arc<dyn DeviceSigner>>,
249    #[cfg(feature = "managed-group-encryption")]
250    pub managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
251}
252
253#[cfg(feature = "managed-group-encryption")]
254#[async_trait]
255pub trait NativeManagedGroupStateStore: Send + Sync {
256    async fn load(
257        &self,
258        app_tag: &str,
259        device_id: &str,
260        avenue_key: &str,
261    ) -> Result<Option<Vec<u8>>>;
262    async fn save(
263        &self,
264        app_tag: &str,
265        device_id: &str,
266        avenue_key: &str,
267        sealed_state: &[u8],
268    ) -> Result<()>;
269}
270
271/// Process-local state store for prototypes and deterministic tests. Shipping
272/// native hosts should provide durable owner-only app-private storage.
273#[cfg(feature = "managed-group-encryption")]
274#[derive(Default)]
275pub struct InMemoryNativeManagedGroupStateStore {
276    states: Mutex<HashMap<String, Vec<u8>>>,
277}
278
279#[cfg(feature = "managed-group-encryption")]
280#[async_trait]
281impl NativeManagedGroupStateStore for InMemoryNativeManagedGroupStateStore {
282    async fn load(
283        &self,
284        app_tag: &str,
285        device_id: &str,
286        avenue_key: &str,
287    ) -> Result<Option<Vec<u8>>> {
288        Ok(self
289            .states
290            .lock()
291            .map_err(|_| anyhow!("managed group state store is poisoned"))?
292            .get(&format!("{app_tag}:{device_id}:{avenue_key}"))
293            .cloned())
294    }
295
296    async fn save(
297        &self,
298        app_tag: &str,
299        device_id: &str,
300        avenue_key: &str,
301        sealed_state: &[u8],
302    ) -> Result<()> {
303        if sealed_state.is_empty() {
304            bail!("managed group sealed state is empty");
305        }
306        self.states
307            .lock()
308            .map_err(|_| anyhow!("managed group state store is poisoned"))?
309            .insert(
310                format!("{app_tag}:{device_id}:{avenue_key}"),
311                sealed_state.to_vec(),
312            );
313        Ok(())
314    }
315}
316
317#[cfg(feature = "managed-group-encryption")]
318#[derive(Debug, Clone, PartialEq, Eq)]
319pub struct NativeManagedRoomPublish {
320    pub message_id: String,
321    pub channel: String,
322    pub priority: u8,
323    pub zone_id: Option<String>,
324    pub payload: Vec<u8>,
325}
326
327#[cfg(feature = "managed-group-encryption")]
328#[derive(Debug, Clone, PartialEq, Eq)]
329pub struct NativeManagedRoomMessage {
330    pub sender_device_id: String,
331    pub message_id: String,
332    pub channel: String,
333    pub priority: u8,
334    pub zone_id: Option<String>,
335    pub payload: Vec<u8>,
336}
337
338/// Signed, application-owned AOI/interest assignment delivered by the room
339/// authority. OpenRTC validates its bounded wire shape and uses it only as a
340/// routing input; the application remains the owner of entity semantics.
341#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
342#[serde(rename_all = "camelCase")]
343pub struct NativeAuthorityInterestAssignment {
344    pub policy_version: String,
345    pub service_id: String,
346    pub generation: u64,
347    pub revision: u64,
348    pub subject_device_id: String,
349    pub shard_id: String,
350    pub priority_device_ids: Vec<String>,
351    pub relevant_entity_ids: Vec<String>,
352    pub issued_at_ms: u64,
353    pub expires_at_ms: u64,
354}
355
356#[derive(Debug, Clone, PartialEq, Eq)]
357pub struct NativeAuthorityAssignment {
358    pub assignment: NativeAuthorityInterestAssignment,
359    pub signature: String,
360}
361
362/// Lightweight, network-idle entrypoint for pure-Rust OpenRTC 2.0 capability
363/// activation. It mirrors the public SDK namespaces without putting consumer
364/// authentication, attestation, or provider-specific types into the runtime.
365pub struct NativeCapabilities {
366    endpoint: String,
367    app_tag: String,
368    device_id: String,
369    platform_type: String,
370    grant_provider: Arc<dyn NativeGatewayGrantProvider>,
371    #[cfg(feature = "managed-group-encryption")]
372    managed_group_signer: Option<Arc<dyn DeviceSigner>>,
373    #[cfg(feature = "managed-group-encryption")]
374    managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
375}
376
377impl NativeCapabilities {
378    /// Construct the production capability namespace from a public developer
379    /// API key. No grant, socket, timer, or provider request occurs here.
380    pub fn new(
381        api_key: &str,
382        device_id: impl Into<String>,
383        platform_type: impl Into<String>,
384        grant_provider: Arc<dyn NativeGatewayGrantProvider>,
385    ) -> Result<Self> {
386        let api_key = crate::validate_api_key(api_key)?;
387        Ok(Self {
388            endpoint: OPENRTC_PRODUCTION_COORDINATION_GATEWAY.to_string(),
389            app_tag: crate::app_tag_from_api_key(api_key),
390            device_id: required("device_id", device_id.into())?,
391            platform_type: required("platform_type", platform_type.into())?,
392            grant_provider,
393            #[cfg(feature = "managed-group-encryption")]
394            managed_group_signer: None,
395            #[cfg(feature = "managed-group-encryption")]
396            managed_group_store: None,
397        })
398    }
399
400    /// Install prompt-free native managed-room persistence. OpenRTC invokes
401    /// the sign-only device key internally and never exports the derived key.
402    #[cfg(feature = "managed-group-encryption")]
403    pub fn with_managed_group_storage(
404        mut self,
405        signer: Arc<dyn DeviceSigner>,
406        store: Arc<dyn NativeManagedGroupStateStore>,
407    ) -> Self {
408        self.managed_group_signer = Some(signer);
409        self.managed_group_store = Some(store);
410        self
411    }
412
413    #[cfg(any(test, feature = "testing-endpoints"))]
414    pub fn with_testing_endpoint(mut self, endpoint: impl Into<String>) -> Result<Self> {
415        self.endpoint = validate_endpoint("endpoint", endpoint.into(), true)?;
416        Ok(self)
417    }
418
419    pub fn devices(&self, principal_id: impl Into<String>) -> Result<NativeCapabilityHandle> {
420        self.open("user", principal_id, None, None)
421    }
422
423    pub fn join_space(&self, id: impl Into<String>) -> Result<NativeCapabilityHandle> {
424        self.open("space", id, None, None)
425    }
426
427    pub fn join_room_with_architecture(
428        &self,
429        id: impl Into<String>,
430        architecture: RoomArchitectureMode,
431    ) -> Result<NativeCapabilityHandle> {
432        self.join_room_with_options(id, architecture, RoomDelivery::Reliable)
433    }
434
435    pub fn join_room_with_options(
436        &self,
437        id: impl Into<String>,
438        architecture: RoomArchitectureMode,
439        delivery: RoomDelivery,
440    ) -> Result<NativeCapabilityHandle> {
441        self.open("room", id, Some(architecture), Some(delivery))
442    }
443
444    pub fn issue_ticket(&self, id: impl Into<String>) -> Result<NativeCapabilityHandle> {
445        self.open("session", id, None, None)
446    }
447
448    fn open(
449        &self,
450        kind: &'static str,
451        id: impl Into<String>,
452        architecture: Option<RoomArchitectureMode>,
453        room_delivery: Option<RoomDelivery>,
454    ) -> Result<NativeCapabilityHandle> {
455        NativeCapabilityHandle::new(GatewayOptions {
456            endpoint: self.endpoint.clone(),
457            app_tag: self.app_tag.clone(),
458            device_id: self.device_id.clone(),
459            platform_type: self.platform_type.clone(),
460            avenue: NativeCoordinationAvenue {
461                kind: kind.to_string(),
462                id: id.into(),
463            },
464            architecture,
465            room_delivery,
466            grant_provider: self.grant_provider.clone(),
467            #[cfg(feature = "managed-group-encryption")]
468            managed_group_signer: self.managed_group_signer.clone(),
469            #[cfg(feature = "managed-group-encryption")]
470            managed_group_store: self.managed_group_store.clone(),
471        })
472    }
473}
474
475/// Disposable OpenRTC 2.0 native avenue.
476///
477/// A handle owns exactly one coordination adapter and therefore exactly one
478/// devices, space, room, or ticket avenue. Constructing the root `Client` does
479/// not create one of these handles and performs no network work. The socket is
480/// opened only when the runtime publishes its initial presence through the
481/// adapter. Closing the handle stops grant refresh, reconnect, and presence
482/// ownership for that avenue without replacing the shared peer runtime.
483pub struct NativeCapabilityHandle {
484    kind: NativeCapabilityKind,
485    id: String,
486    signaling: Arc<NativeCoordinationGatewaySignaling>,
487    closed: Arc<AtomicBool>,
488}
489
490pub(crate) struct NativeCapabilityCloser {
491    signaling: Arc<NativeCoordinationGatewaySignaling>,
492    closed: Arc<AtomicBool>,
493}
494
495#[derive(Debug, Clone, Copy, PartialEq, Eq)]
496pub enum NativeCapabilityKind {
497    Devices,
498    Space,
499    Room,
500    Ticket,
501}
502
503impl NativeCapabilityKind {
504    fn from_avenue_kind(value: &str) -> Result<Self> {
505        match value {
506            "user" => Ok(Self::Devices),
507            "space" => Ok(Self::Space),
508            "room" => Ok(Self::Room),
509            "session" => Ok(Self::Ticket),
510            _ => bail!("native coordination avenue kind is invalid"),
511        }
512    }
513}
514
515impl NativeCapabilityHandle {
516    /// Create one inactive native capability handle. This allocates the local
517    /// actor only; it does not request a grant or connect to the gateway.
518    pub fn new(options: GatewayOptions) -> Result<Self> {
519        let kind = NativeCapabilityKind::from_avenue_kind(&options.avenue.kind)?;
520        let id = required("avenue.id", options.avenue.id.clone())?;
521        let signaling = NativeCoordinationGatewaySignaling::new(options)?;
522        Ok(Self {
523            kind,
524            id,
525            signaling,
526            closed: Arc::new(AtomicBool::new(false)),
527        })
528    }
529
530    pub fn kind(&self) -> NativeCapabilityKind {
531        self.kind
532    }
533
534    pub fn id(&self) -> &str {
535        &self.id
536    }
537
538    pub fn signaling(&self) -> Arc<dyn SignalingBackend> {
539        self.signaling.clone()
540    }
541
542    pub(crate) async fn compose_authority_client(
543        &self,
544        builder: crate::client::ClientBuilder,
545    ) -> Result<Arc<crate::Client>> {
546        if self.is_closed()
547            || self.kind != NativeCapabilityKind::Room
548            || self.signaling.architecture != Some(RoomArchitectureMode::Authority)
549        {
550            bail!("native authority composition requires an open authority room");
551        }
552        let mut binding = self.signaling.authority_client.lock().await;
553        // Close can win while composition is waiting for this owner lock.
554        if self.is_closed() {
555            bail!("native authority is closed");
556        }
557        if binding.is_some() {
558            bail!("native authority already has a peer-session owner");
559        }
560        let client = Arc::new(builder.signaling_backend(self.signaling()).build());
561        if client.app_tag() != self.signaling.app_tag {
562            bail!("native authority client app identity does not match its capability");
563        }
564        let scope = format!("v2:room:{}", self.id);
565        client
566            .session_token_registry
567            .require_scope_peer_admission(&scope)
568            .map_err(anyhow::Error::msg)?;
569        client.ensure_default_admission_gate(&scope);
570        client
571            .start_external_auto_connect(
572                format!("{}:native-root", self.signaling.app_tag),
573                self.signaling.device_id.clone(),
574            )
575            .await?;
576        *binding = Some(AuthorityPeerBinding {
577            client: Arc::downgrade(&client),
578            scope,
579            revision: 0,
580            admission_revision: 0,
581            last_payload: String::new(),
582            retiring_tickets: Vec::new(),
583            closed: false,
584        });
585        Ok(client)
586    }
587
588    pub fn room_architecture(&self) -> Result<Option<RoomArchitectureSnapshot>> {
589        if self.kind != NativeCapabilityKind::Room {
590            bail!("native room architecture is available only for room capabilities");
591        }
592        Ok(self.signaling.room_architecture())
593    }
594
595    pub fn subscribe_room_architecture(
596        &self,
597    ) -> Result<tokio::sync::watch::Receiver<Option<RoomArchitectureSnapshot>>> {
598        if self.kind != NativeCapabilityKind::Room {
599            bail!("native room architecture is available only for room capabilities");
600        }
601        Ok(self.signaling.subscribe_room_architecture())
602    }
603
604    pub async fn publish_authority_assignment(
605        &self,
606        assignment: NativeAuthorityInterestAssignment,
607        signature: impl Into<String>,
608    ) -> Result<()> {
609        if self.kind != NativeCapabilityKind::Room {
610            bail!("authority assignments are available only for room capabilities");
611        }
612        let signature = signature.into();
613        validate_authority_assignment(&assignment, &signature)?;
614        self.signaling
615            .publish_authority_assignment(assignment, signature)
616            .await
617    }
618
619    pub fn authority_assignment(&self) -> Result<Option<NativeAuthorityAssignment>> {
620        if self.kind != NativeCapabilityKind::Room {
621            bail!("authority assignments are available only for room capabilities");
622        }
623        Ok(self.signaling.authority_assignment())
624    }
625
626    pub fn subscribe_authority_assignments(
627        &self,
628    ) -> Result<watch::Receiver<Option<NativeAuthorityAssignment>>> {
629        if self.kind != NativeCapabilityKind::Room {
630            bail!("authority assignments are available only for room capabilities");
631        }
632        Ok(self.signaling.subscribe_authority_assignments())
633    }
634
635    pub fn subscribe_service_errors(
636        &self,
637    ) -> Result<broadcast::Receiver<NativeServiceErrorObservation>> {
638        if self.is_closed() {
639            bail!("native capability is closed");
640        }
641        Ok(self.signaling.subscribe_service_errors())
642    }
643
644    #[cfg(feature = "managed-group-encryption")]
645    pub async fn publish_managed_room(&self, batch: Vec<NativeManagedRoomPublish>) -> Result<()> {
646        if self.kind != NativeCapabilityKind::Room {
647            bail!("managed fanout is available only for room capabilities");
648        }
649        self.signaling.publish_managed_room(batch).await
650    }
651
652    #[cfg(feature = "managed-group-encryption")]
653    pub fn subscribe_managed_room(&self) -> Result<broadcast::Receiver<NativeManagedRoomMessage>> {
654        if self.kind != NativeCapabilityKind::Room {
655            bail!("managed fanout is available only for room capabilities");
656        }
657        Ok(self.signaling.subscribe_managed_room())
658    }
659
660    pub fn is_closed(&self) -> bool {
661        self.closed.load(Ordering::Acquire)
662    }
663
664    pub async fn close(&self) {
665        if !self.closed.swap(true, Ordering::AcqRel) {
666            self.signaling.stop().await;
667        }
668    }
669
670    pub(crate) fn closer(&self) -> NativeCapabilityCloser {
671        NativeCapabilityCloser {
672            signaling: self.signaling.clone(),
673            closed: self.closed.clone(),
674        }
675    }
676}
677
678pub(crate) fn validate_native_authority_assignment_targets(
679    priority_device_ids: &[String],
680    relevant_entity_ids: &[String],
681) -> Result<()> {
682    let valid_id = |value: &str, max: usize| {
683        !value.is_empty()
684            && value.len() <= max
685            && value.bytes().all(|byte| {
686                byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'.' | b':' | b'@' | b'-')
687            })
688    };
689    if priority_device_ids.len() > MAX_AUTHORITY_PRIORITY_PEERS
690        || relevant_entity_ids.len() > MAX_AUTHORITY_RELEVANT_ENTITIES
691        || priority_device_ids
692            .iter()
693            .any(|value| !valid_id(value, 160))
694        || relevant_entity_ids
695            .iter()
696            .any(|value| !valid_id(value, 160))
697        || priority_device_ids
698            .iter()
699            .collect::<std::collections::HashSet<_>>()
700            .len()
701            != priority_device_ids.len()
702        || relevant_entity_ids
703            .iter()
704            .collect::<std::collections::HashSet<_>>()
705            .len()
706            != relevant_entity_ids.len()
707    {
708        bail!("native authority assignment targets are invalid");
709    }
710    Ok(())
711}
712
713fn validate_authority_assignment(
714    assignment: &NativeAuthorityInterestAssignment,
715    signature: &str,
716) -> Result<()> {
717    validate_native_authority_assignment_targets(
718        &assignment.priority_device_ids,
719        &assignment.relevant_entity_ids,
720    )?;
721    let valid_id = |value: &str, max: usize| {
722        !value.is_empty()
723            && value.len() <= max
724            && value.bytes().all(|byte| {
725                byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'.' | b':' | b'@' | b'-')
726            })
727    };
728    if assignment.policy_version != "room-authority-assignment-v1"
729        || !valid_id(&assignment.service_id, 80)
730        || !valid_id(&assignment.subject_device_id, 160)
731        || !valid_id(&assignment.shard_id, 80)
732        || assignment.generation == 0
733        || assignment.revision == 0
734        || assignment.expires_at_ms <= assignment.issued_at_ms
735        || assignment.issued_at_ms > now_ms().saturating_add(30_000)
736        || assignment.expires_at_ms <= now_ms()
737        || assignment
738            .expires_at_ms
739            .saturating_sub(assignment.issued_at_ms)
740            > 2 * 60_000
741    {
742        bail!("native authority assignment is invalid");
743    }
744    let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
745        .decode(signature)
746        .context("decode native authority assignment signature")?;
747    if decoded.len() != 64 {
748        bail!("native authority assignment signature is invalid");
749    }
750    Ok(())
751}
752
753impl NativeCapabilityCloser {
754    pub(crate) async fn close(&self) {
755        if !self.closed.swap(true, Ordering::AcqRel) {
756            self.signaling.stop().await;
757        }
758    }
759}
760
761impl Drop for NativeCapabilityHandle {
762    fn drop(&mut self) {
763        if !self.closed.swap(true, Ordering::AcqRel) {
764            let _ = self.signaling.commands.try_send(Command::Stop);
765        }
766    }
767}
768
769#[derive(Debug, Clone, PartialEq, Eq)]
770struct DesiredPresence {
771    user_id: String,
772    local_node_id: String,
773    ticket: String,
774    device_name: String,
775    metadata: Option<String>,
776    ttl_ms: u64,
777    online: bool,
778}
779
780#[derive(Debug, Clone, Default, Serialize)]
781#[serde(rename_all = "camelCase")]
782struct DevicePatch {
783    #[serde(skip_serializing_if = "Option::is_none")]
784    device_name: Option<String>,
785    #[serde(skip_serializing_if = "Option::is_none")]
786    capabilities: Option<DeviceCapabilities>,
787    #[serde(skip_serializing_if = "Option::is_none")]
788    metadata: Option<String>,
789    #[serde(skip_serializing_if = "Option::is_none")]
790    excluded_peers: Option<Vec<String>>,
791}
792
793#[derive(Debug, Serialize)]
794#[serde(rename_all = "camelCase")]
795struct OutboundGatewayDevice {
796    device_id: String,
797    runtime_instance_id: String,
798    node_id: String,
799    device_name: String,
800    platform_type: String,
801    ticket: String,
802    #[serde(skip_serializing_if = "Option::is_none")]
803    metadata: Option<String>,
804    #[serde(skip_serializing_if = "Option::is_none")]
805    capabilities: Option<DeviceCapabilities>,
806    excluded_peers: Vec<String>,
807    online: bool,
808}
809
810#[derive(Debug)]
811enum Command {
812    Publish {
813        desired: DesiredPresence,
814        reply: oneshot::Sender<Result<()>>,
815    },
816    Patch {
817        patch: DevicePatch,
818        reply: oneshot::Sender<Result<()>>,
819    },
820    Offline {
821        reply: oneshot::Sender<Result<()>>,
822    },
823    Delete {
824        user_id: String,
825        device_id: String,
826        reply: oneshot::Sender<Result<()>>,
827    },
828    SendSignal {
829        target_device_id: String,
830        payload: String,
831        state: Option<String>,
832        reply_payload: Option<String>,
833        reply: oneshot::Sender<Result<String>>,
834    },
835    PutSession {
836        session_id: String,
837        session: serde_json::Value,
838        expires_at_ms: i64,
839        reply: oneshot::Sender<Result<()>>,
840    },
841    PublishAuthorityAssignment {
842        assignment: NativeAuthorityInterestAssignment,
843        signature: String,
844        reply: oneshot::Sender<Result<()>>,
845    },
846    #[cfg(feature = "managed-group-encryption")]
847    PublishManaged {
848        batch: Vec<NativeManagedRoomPublish>,
849        reply: oneshot::Sender<Result<()>>,
850    },
851    Stop,
852}
853
854#[derive(Debug, Clone, Deserialize)]
855#[serde(rename_all = "camelCase")]
856struct GatewayDevice {
857    #[serde(default)]
858    user_id: Option<String>,
859    device_id: String,
860    runtime_instance_id: String,
861    node_id: String,
862    device_name: String,
863    platform_type: String,
864    ticket: String,
865    #[serde(default)]
866    metadata: Option<String>,
867    #[serde(default)]
868    capabilities: Option<DeviceCapabilities>,
869    #[serde(default)]
870    excluded_peers: Vec<String>,
871    online: bool,
872    updated_at_ms: i64,
873    expires_at_ms: i64,
874}
875
876#[derive(Debug, Clone, Deserialize)]
877#[serde(rename_all = "camelCase")]
878struct GatewayMember {
879    #[serde(default)]
880    user_id: Option<String>,
881    device_id: String,
882    device_name: String,
883    platform_type: String,
884    #[serde(default)]
885    metadata: Option<String>,
886    #[serde(default)]
887    capabilities: Option<DeviceCapabilities>,
888    online: bool,
889    updated_at_ms: i64,
890    expires_at_ms: i64,
891}
892
893#[derive(Debug, Clone, Deserialize)]
894#[serde(rename_all = "camelCase")]
895struct GatewayArchitectureLease {
896    policy_version: String,
897    epoch: u64,
898    #[serde(default)]
899    previous_architecture_epoch: Option<u64>,
900    requested_mode: RoomArchitectureMode,
901    effective_mode: EffectiveRoomArchitecture,
902    phase: RoomArchitecturePhase,
903    reason: RoomArchitectureReason,
904    held_credits_microusd: u64,
905    quote_expires_at_ms: u64,
906    #[serde(default)]
907    reservation_id: Option<String>,
908}
909
910#[derive(Debug, Clone, Deserialize)]
911#[serde(rename_all = "camelCase")]
912struct GatewayGroupEncryptionLease {
913    policy_version: String,
914    epoch: u64,
915    #[serde(default)]
916    previous_encryption_epoch: Option<u64>,
917    phase: String,
918    committer_device_id: String,
919    group_id_hash: String,
920}
921
922#[cfg(feature = "managed-group-encryption")]
923#[derive(Debug, Clone, Serialize, Deserialize)]
924#[serde(rename_all = "camelCase")]
925struct NativeManagedEncryptedEnvelope {
926    message_id: String,
927    channel: String,
928    priority: u8,
929    scope: NativeManagedScope,
930    ciphertext: String,
931}
932
933#[cfg(feature = "managed-group-encryption")]
934#[derive(Debug, Clone, Serialize, Deserialize)]
935#[serde(tag = "kind", rename_all = "lowercase")]
936enum NativeManagedScope {
937    Global,
938    Zone {
939        #[serde(rename = "zoneId")]
940        zone_id: String,
941    },
942}
943
944#[derive(Debug, Clone, Deserialize)]
945#[serde(rename_all = "camelCase")]
946struct GatewayTopologyLease {
947    schema_version: u8,
948    topology_revision: u64,
949    #[serde(default)]
950    previous_topology_revision: Option<u64>,
951    avenue: NativeCoordinationAvenue,
952    grant_jti: String,
953    expires_at_ms: u64,
954    architecture: GatewayArchitectureLease,
955    #[serde(default)]
956    group_encryption: Option<GatewayGroupEncryptionLease>,
957    #[serde(default)]
958    admission_peers: Option<Vec<GatewayAdmissionPeer>>,
959    active: Vec<GatewayDevice>,
960    backups: Vec<GatewayDevice>,
961}
962
963#[derive(Debug, Clone, Deserialize)]
964#[serde(rename_all = "camelCase")]
965struct GatewayAdmissionPeer {
966    device_id: String,
967    node_id: String,
968}
969
970#[derive(Debug, Deserialize)]
971#[serde(tag = "type")]
972enum ServerFrame {
973    #[serde(rename = "ready")]
974    Ready {
975        #[serde(rename = "budgetRemainingMicrousd")]
976        _budget_remaining_microusd: i64,
977        #[serde(default, rename = "leaseRefreshMode")]
978        lease_refresh_mode: Option<String>,
979    },
980    #[serde(rename = "auth.refreshed")]
981    AuthRefreshed {
982        #[serde(rename = "expiresAtMs")]
983        _expires_at_ms: u64,
984    },
985    #[serde(rename = "lease.refreshed")]
986    LeaseRefreshed {
987        #[serde(rename = "idempotencyKey")]
988        idempotency_key: String,
989        #[serde(rename = "expiresAtMs")]
990        expires_at_ms: u64,
991    },
992    #[serde(rename = "roster.snapshot")]
993    RosterSnapshot { devices: Vec<GatewayDevice> },
994    #[serde(rename = "membership.snapshot")]
995    MembershipSnapshot { members: Vec<GatewayMember> },
996    #[serde(rename = "membership.page")]
997    MembershipPage {
998        #[serde(rename = "snapshotId")]
999        snapshot_id: String,
1000        #[serde(rename = "pageIndex")]
1001        page_index: usize,
1002        #[serde(rename = "pageCount")]
1003        page_count: usize,
1004        #[serde(rename = "memberCount")]
1005        member_count: usize,
1006        members: Vec<GatewayMember>,
1007    },
1008    #[serde(rename = "membership.changed")]
1009    MembershipChanged {
1010        operation: String,
1011        member: GatewayMember,
1012    },
1013    #[serde(rename = "topology.lease")]
1014    TopologyLease { lease: GatewayTopologyLease },
1015    #[cfg(feature = "managed-group-encryption")]
1016    #[serde(rename = "managed.prepare.page")]
1017    ManagedPreparePage {
1018        #[serde(flatten)]
1019        page: ManagedPreparePage,
1020    },
1021    #[cfg(feature = "managed-group-encryption")]
1022    #[serde(rename = "managed.commit.chunk.received")]
1023    ManagedArtifactChunk {
1024        #[serde(flatten)]
1025        chunk: ManagedArtifactChunk,
1026    },
1027    #[cfg(feature = "managed-group-encryption")]
1028    #[serde(rename = "managed.received")]
1029    ManagedReceived {
1030        #[serde(rename = "senderDeviceId")]
1031        sender_device_id: String,
1032        #[serde(rename = "architectureEpoch")]
1033        architecture_epoch: u64,
1034        #[serde(rename = "encryptionEpoch")]
1035        encryption_epoch: u64,
1036        batch: Vec<NativeManagedEncryptedEnvelope>,
1037    },
1038    #[serde(rename = "authority.assignment")]
1039    AuthorityAssignment {
1040        assignment: NativeAuthorityInterestAssignment,
1041        signature: String,
1042    },
1043    #[serde(rename = "presence.changed")]
1044    PresenceChanged {
1045        operation: String,
1046        device: GatewayDevice,
1047    },
1048    #[serde(rename = "session.changed")]
1049    SessionChanged {
1050        operation: String,
1051        #[serde(rename = "sessionId")]
1052        session_id: String,
1053        #[serde(default)]
1054        session: Option<serde_json::Value>,
1055    },
1056    #[serde(rename = "ack")]
1057    Ack {
1058        #[serde(rename = "idempotencyKey")]
1059        idempotency_key: String,
1060    },
1061    #[serde(rename = "error")]
1062    Error {
1063        code: String,
1064        message: String,
1065        #[serde(default, rename = "idempotencyKey")]
1066        idempotency_key: Option<String>,
1067        retryable: bool,
1068        #[serde(flatten)]
1069        metadata: serde_json::Map<String, serde_json::Value>,
1070    },
1071    #[serde(rename = "signal.received")]
1072    SignalReceived {
1073        #[serde(rename = "signalId")]
1074        _signal_id: String,
1075        #[serde(rename = "senderDeviceId")]
1076        sender_device_id: String,
1077        payload: String,
1078        #[serde(default)]
1079        state: Option<String>,
1080        #[serde(default, rename = "replyPayload")]
1081        reply_payload: Option<String>,
1082        #[serde(rename = "createdAtMs")]
1083        _created_at_ms: i64,
1084    },
1085    #[serde(rename = "pong")]
1086    Pong,
1087}
1088
1089struct SharedState {
1090    desired: RwLock<Option<DesiredPresence>>,
1091    applied_presence: RwLock<Option<DesiredPresence>>,
1092    staged_patch: RwLock<DevicePatch>,
1093    members: RwLock<HashMap<String, GatewayMember>>,
1094    membership_pages: Mutex<Option<PendingMembershipPages>>,
1095    // Preserve active/backup intent from the gateway; a public roster is not
1096    // permission to dial, and a backup must not become an active authority edge.
1097    topology_routes: RwLock<HashMap<String, (GatewayDevice, bool)>>,
1098    topology_revision: RwLock<u64>,
1099    encryption_epoch: RwLock<u64>,
1100    active_grant: RwLock<Option<(String, u64)>>,
1101    devices: RwLock<HashMap<String, Device>>,
1102    device_events: broadcast::Sender<Vec<DeviceEvent>>,
1103    session_events: broadcast::Sender<Vec<SessionEvent>>,
1104    service_errors: broadcast::Sender<NativeServiceErrorObservation>,
1105    room_architecture: watch::Sender<Option<RoomArchitectureSnapshot>>,
1106    authority_assignment: watch::Sender<Option<NativeAuthorityAssignment>>,
1107    pending_messages: Mutex<Vec<SignalingEnvelope>>,
1108    #[cfg(feature = "managed-group-encryption")]
1109    managed_group: tokio::sync::Mutex<Option<ManagedGroupController>>,
1110    #[cfg(feature = "managed-group-encryption")]
1111    managed_messages: broadcast::Sender<NativeManagedRoomMessage>,
1112}
1113
1114struct PendingMembershipPages {
1115    snapshot_id: String,
1116    next_page_index: usize,
1117    page_count: usize,
1118    member_count: usize,
1119    members: HashMap<String, GatewayMember>,
1120}
1121
1122// Subordinate input binding, not a second peer/session registry. The gateway
1123// actor serializes changes; Rust Client owns dialing, admission and recovery.
1124// Reviewed 2026-09-04; consumers: RoomAuthority and the bundled Node host.
1125struct AuthorityPeerBinding {
1126    client: std::sync::Weak<crate::Client>,
1127    scope: String,
1128    revision: u64,
1129    // Unlike per-socket topology revision, survives gateway reconnects.
1130    admission_revision: u64,
1131    last_payload: String,
1132    // At most one normal overlap: renewal is one minute before a 15-minute
1133    // ticket expires. Old admissions are revoked by the existing Client owner.
1134    retiring_tickets: Vec<(String, u64)>,
1135    closed: bool,
1136}
1137
1138pub(crate) struct NativeCoordinationGatewaySignaling {
1139    endpoint: String,
1140    app_tag: String,
1141    device_id: String,
1142    avenue: NativeCoordinationAvenue,
1143    architecture: Option<RoomArchitectureMode>,
1144    room_delivery: Option<RoomDelivery>,
1145    runtime_instance_id: String,
1146    platform_type: String,
1147    grant_provider: Arc<dyn NativeGatewayGrantProvider>,
1148    shared: Arc<SharedState>,
1149    commands: mpsc::Sender<Command>,
1150    authority_client: tokio::sync::Mutex<Option<AuthorityPeerBinding>>,
1151    #[cfg(feature = "managed-group-encryption")]
1152    managed_group_signer: Option<Arc<dyn DeviceSigner>>,
1153    #[cfg(feature = "managed-group-encryption")]
1154    managed_group_store: Option<Arc<dyn NativeManagedGroupStateStore>>,
1155}
1156
1157impl NativeCoordinationGatewaySignaling {
1158    pub(crate) fn new(options: GatewayOptions) -> Result<Arc<Self>> {
1159        let endpoint = validate_endpoint("endpoint", options.endpoint, true)?;
1160        let app_tag = required("app_tag", options.app_tag)?;
1161        let device_id = required("device_id", options.device_id)?;
1162        let platform_type = required("platform_type", options.platform_type)?;
1163        let avenue = options.avenue;
1164        if !matches!(avenue.kind.as_str(), "user" | "space" | "room" | "session") {
1165            bail!("native coordination avenue kind is invalid");
1166        }
1167        let avenue_id = required("avenue.id", avenue.id)?;
1168        if options.architecture.is_some() && avenue.kind != "room" {
1169            bail!("native room architecture is valid only for room avenues");
1170        }
1171        if options.room_delivery.is_some() && avenue.kind != "room" {
1172            bail!("native room delivery is valid only for room avenues");
1173        }
1174        let (device_events, _) = broadcast::channel(64);
1175        let (session_events, _) = broadcast::channel(64);
1176        let (service_errors, _) = broadcast::channel(64);
1177        let (room_architecture, _) = watch::channel(None);
1178        let (authority_assignment, _) = watch::channel(None);
1179        #[cfg(feature = "managed-group-encryption")]
1180        let (managed_messages, _) = broadcast::channel(256);
1181        let shared = Arc::new(SharedState {
1182            desired: RwLock::new(None),
1183            applied_presence: RwLock::new(None),
1184            staged_patch: RwLock::new(DevicePatch::default()),
1185            members: RwLock::new(HashMap::new()),
1186            membership_pages: Mutex::new(None),
1187            topology_routes: RwLock::new(HashMap::new()),
1188            topology_revision: RwLock::new(0),
1189            encryption_epoch: RwLock::new(0),
1190            active_grant: RwLock::new(None),
1191            devices: RwLock::new(HashMap::new()),
1192            device_events,
1193            session_events,
1194            service_errors,
1195            room_architecture,
1196            authority_assignment,
1197            pending_messages: Mutex::new(Vec::new()),
1198            #[cfg(feature = "managed-group-encryption")]
1199            managed_group: tokio::sync::Mutex::new(None),
1200            #[cfg(feature = "managed-group-encryption")]
1201            managed_messages,
1202        });
1203        let (commands, receiver) = mpsc::channel(64);
1204        let adapter = Arc::new(Self {
1205            endpoint,
1206            app_tag,
1207            device_id,
1208            avenue: NativeCoordinationAvenue {
1209                kind: avenue.kind,
1210                id: avenue_id,
1211            },
1212            architecture: options.architecture,
1213            room_delivery: options.room_delivery,
1214            runtime_instance_id: format!("runtime:{}", Uuid::new_v4()),
1215            platform_type,
1216            grant_provider: options.grant_provider,
1217            shared,
1218            commands,
1219            authority_client: tokio::sync::Mutex::new(None),
1220            #[cfg(feature = "managed-group-encryption")]
1221            managed_group_signer: options.managed_group_signer,
1222            #[cfg(feature = "managed-group-encryption")]
1223            managed_group_store: options.managed_group_store,
1224        });
1225        tokio::spawn(run_actor(adapter.clone(), receiver));
1226        Ok(adapter)
1227    }
1228
1229    pub async fn stop(&self) {
1230        self.stop_authority_client().await;
1231        let _ = self.commands.send(Command::Stop).await;
1232    }
1233
1234    async fn authority_ticket_maintenance_delay(&self) -> Duration {
1235        let binding = self.authority_client.lock().await;
1236        let Some(binding) = binding.as_ref().filter(|binding| !binding.closed && binding.client.strong_count() > 0) else {
1237            return Duration::from_secs(365 * 24 * 60 * 60);
1238        };
1239        let desired = self.shared.desired.read().await;
1240        let renewal = desired.as_ref().map(|desired| {
1241            let (_, suffix) = crate::session_token::split_ticket(&desired.ticket);
1242            suffix
1243                .and_then(crate::session_token::decode_token_payload)
1244                .and_then(|payload| payload.expires_at_ms)
1245                .map(|expires| expires.saturating_sub(AUTHORITY_TICKET_RENEWAL_SKEW_MS))
1246                .unwrap_or(0)
1247        });
1248        let deadline = binding
1249            .retiring_tickets
1250            .iter()
1251            .map(|(_, expiry)| *expiry)
1252            .chain(renewal)
1253            .min();
1254        deadline.map_or(Duration::from_secs(365 * 24 * 60 * 60), |deadline| {
1255            Duration::from_millis(deadline.saturating_sub(now_ms()).max(1))
1256        })
1257    }
1258
1259    /// Maintenance of the locally authored service ticket, never a peer's
1260    /// advertised credential. Runs in the existing gateway actor before grant
1261    /// minting or at its renewal deadline; it performs no hosted work itself.
1262    async fn maintain_authority_ticket(&self) -> Result<bool> {
1263        self.maintain_authority_publication(None).await
1264    }
1265
1266    async fn maintain_authority_publication(&self, pending: Option<&mut DesiredPresence>) -> Result<bool> {
1267        self.expire_retiring_authority_tickets().await;
1268        let mut binding = self.authority_client.lock().await;
1269        let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
1270            return Ok(false);
1271        };
1272        let Some(client) = binding.client.upgrade() else {
1273            return Ok(false);
1274        };
1275        let now = now_ms();
1276        let mut desired = self.shared.desired.write().await;
1277        let Some(desired) = desired.as_mut() else {
1278            return Ok(false);
1279        };
1280        let (ticket, suffix) = crate::session_token::split_ticket(&desired.ticket);
1281        // Decode the locally stored ticket even after expiry (e.g. sleep).
1282        // Remote ingress continues to use the expiry-validating decoder.
1283        let payload = suffix
1284            .and_then(crate::session_token::decode_token_payload)
1285            .ok_or_else(|| anyhow!("authority presence requires a scoped expiring ticket"))?;
1286        if payload.scope.as_str() != binding.scope
1287            || payload.ticket_hash.as_deref()
1288                != Some(crate::session_token::endpoint_ticket_hash(ticket).as_str())
1289        {
1290            bail!("authority presence ticket scope or endpoint binding is invalid");
1291        }
1292        let expires = payload
1293            .expires_at_ms
1294            .ok_or_else(|| anyhow!("authority ticket must expire"))?;
1295        if expires.saturating_sub(now) > AUTHORITY_TICKET_RENEWAL_SKEW_MS {
1296            return Ok(false);
1297        }
1298        // Refuse unexpected repeated rotation rather than accumulate bearer
1299        // state during a failed publication or a caller-driven rewrite.
1300        if binding.retiring_tickets.len() >= 2 {
1301            bail!("authority ticket overlap capacity exhausted");
1302        }
1303        let replacement = client
1304            .endpoint_ticket_with_token(&binding.scope, payload.max_connections)
1305            .await?;
1306        // A local rotation is not a caller superseding its pending publication.
1307        // Update the matching acknowledgement snapshot under the same lock.
1308        if let Some(pending) = pending.filter(|pending| *pending == desired) {
1309            pending.ticket = replacement.clone();
1310        }
1311        desired.ticket = replacement;
1312        if expires <= now {
1313            client.revoke_session_token(&payload.token);
1314        } else {
1315            binding.retiring_tickets.push((payload.token, expires));
1316        }
1317        Ok(true)
1318    }
1319
1320    // Serviced by the existing actor while its one network attempt is pending.
1321    // Do not rotate the fingerprint-bound ticket beneath an in-flight grant.
1322    async fn expire_retiring_authority_tickets(&self) -> Duration {
1323        let mut binding = self.authority_client.lock().await;
1324        let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
1325            return Duration::from_secs(365 * 24 * 60 * 60);
1326        };
1327        let Some(client) = binding.client.upgrade() else {
1328            return Duration::from_secs(365 * 24 * 60 * 60);
1329        };
1330        let now = now_ms();
1331        binding.retiring_tickets.retain(|(token, expiry)| {
1332            if *expiry > now { return true; }
1333            client.revoke_session_token(token);
1334            false
1335        });
1336        binding.retiring_tickets.iter().map(|(_, expiry)| *expiry).min()
1337            .map_or(Duration::from_secs(365 * 24 * 60 * 60), |expiry| {
1338                Duration::from_millis(expiry.saturating_sub(now).max(1))
1339            })
1340    }
1341
1342    async fn forward_authority_peers(&self) -> Result<()> {
1343        let mut binding = self.authority_client.lock().await;
1344        let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
1345            return Ok(());
1346        };
1347        let Some(client) = binding.client.upgrade() else {
1348            return Ok(());
1349        };
1350        let mut peers = Vec::new();
1351        if self.room_architecture().is_some_and(|snapshot| {
1352            snapshot.effective == EffectiveRoomArchitecture::Authority
1353                && snapshot.phase == RoomArchitecturePhase::Settled
1354        }) {
1355            let members = self.shared.members.read().await;
1356            let routes = self.shared.topology_routes.read().await;
1357            for (device, active) in routes.values() {
1358                if !active
1359                    || !device.online
1360                    || device.device_id == self.device_id
1361                    || !members
1362                        .get(&device.device_id)
1363                        .is_some_and(|member| member.online)
1364                {
1365                    continue;
1366                }
1367                peers.push(serde_json::json!({
1368                    "deviceId": device.device_id,
1369                    "nodeId": device.node_id,
1370                    "ticket": device.ticket,
1371                    "online": true,
1372                    "sessionId": device.runtime_instance_id,
1373                    "excludedPeers": device.excluded_peers,
1374                }));
1375            }
1376        }
1377        peers.sort_by(|left, right| left["deviceId"].as_str().cmp(&right["deviceId"].as_str()));
1378        let payload = serde_json::to_string(&peers)?;
1379        if payload != binding.last_payload {
1380            let revision = binding
1381                .revision
1382                .checked_add(1)
1383                .ok_or_else(|| anyhow!("native authority input revision exhausted"))?;
1384            if !client
1385                .submit_external_desired_peers(revision, &payload)
1386                .await?
1387            {
1388                bail!("native authority peer input rejected by its runtime owner");
1389            }
1390            binding.revision = revision;
1391            binding.last_payload = payload;
1392        }
1393        Ok(())
1394    }
1395
1396    async fn handle_gateway_failure(&self, error: &anyhow::Error) -> bool {
1397        if let Some(service_error) = gateway_service_error(error) {
1398            let _ = self.shared.service_errors.send(NativeServiceErrorObservation {
1399                avenue: self.avenue.clone(),
1400                runtime_instance_id: self.runtime_instance_id.clone(),
1401                service_error: service_error.clone(),
1402            });
1403        }
1404        // Reviewed 2026-09-05; gateway actor -> existing room admission owner.
1405        // A revoked credential is not a network outage. Withdraw permissions
1406        // before returning the terminal verdict, including held stream access.
1407        // Ordinary network/budget failures retain their independent lease bounds.
1408        if is_gateway_revocation(error) {
1409            self.stop_authority_client().await;
1410        }
1411        is_retryable_gateway_error(error)
1412    }
1413
1414    fn subscribe_service_errors(&self) -> broadcast::Receiver<NativeServiceErrorObservation> {
1415        self.shared.service_errors.subscribe()
1416    }
1417
1418    async fn stop_authority_client(&self) {
1419        let mut binding = self.authority_client.lock().await;
1420        let Some(binding) = binding.as_mut().filter(|binding| !binding.closed) else {
1421            return;
1422        };
1423        binding.closed = true;
1424        if let Some(client) = binding.client.upgrade() {
1425            binding.admission_revision += 1;
1426            if let Err(error) = client.session_token_registry.update_scope_peer_admission(
1427                &binding.scope,
1428                binding.admission_revision,
1429                0,
1430                Vec::new(),
1431            ) {
1432                eprintln!("[openrtc][authority-close] admission withdrawal failed: {error}");
1433            }
1434            // Keep tokenless incoming traffic blocked throughout scope revocation,
1435            // even when callers retain the Client/endpoint after closing the room.
1436            client.register_session_token(
1437                crate::session_token::generate_token(),
1438                "authority-closed".to_string(),
1439                0,
1440            );
1441            let revoked = client.begin_revoke_tokens_by_scope(&binding.scope);
1442            if let Some(revision) = binding.revision.checked_add(1) {
1443                if let Err(error) = client.submit_external_desired_peers(revision, "[]").await {
1444                    eprintln!("[openrtc][authority-close] peer withdrawal failed: {error:#}");
1445                }
1446            }
1447            client.stop_external_auto_connect().await;
1448            client
1449                .finish_revoke_tokens_by_scope(&binding.scope, &revoked)
1450                .await;
1451        }
1452    }
1453
1454    fn room_architecture(&self) -> Option<RoomArchitectureSnapshot> {
1455        self.shared.room_architecture.borrow().clone()
1456    }
1457
1458    fn subscribe_room_architecture(&self) -> watch::Receiver<Option<RoomArchitectureSnapshot>> {
1459        self.shared.room_architecture.subscribe()
1460    }
1461
1462    async fn publish_authority_assignment(
1463        &self,
1464        assignment: NativeAuthorityInterestAssignment,
1465        signature: String,
1466    ) -> Result<()> {
1467        validate_authority_assignment(&assignment, &signature)?;
1468        self.request_unit(|reply| Command::PublishAuthorityAssignment {
1469            assignment,
1470            signature,
1471            reply,
1472        })
1473        .await
1474    }
1475
1476    fn authority_assignment(&self) -> Option<NativeAuthorityAssignment> {
1477        self.shared
1478            .authority_assignment
1479            .borrow()
1480            .clone()
1481            .filter(|value| value.assignment.expires_at_ms > now_ms())
1482    }
1483
1484    fn subscribe_authority_assignments(
1485        &self,
1486    ) -> watch::Receiver<Option<NativeAuthorityAssignment>> {
1487        self.shared.authority_assignment.subscribe()
1488    }
1489
1490    #[cfg(feature = "managed-group-encryption")]
1491    async fn publish_managed_room(&self, batch: Vec<NativeManagedRoomPublish>) -> Result<()> {
1492        if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
1493            bail!(
1494                "managed room publish batch must contain 1..={MAX_MANAGED_BATCH_MESSAGES} messages"
1495            );
1496        }
1497        self.request_unit(|reply| Command::PublishManaged { batch, reply })
1498            .await
1499    }
1500
1501    #[cfg(feature = "managed-group-encryption")]
1502    fn subscribe_managed_room(&self) -> broadcast::Receiver<NativeManagedRoomMessage> {
1503        self.shared.managed_messages.subscribe()
1504    }
1505
1506    #[cfg(feature = "managed-group-encryption")]
1507    fn avenue_key(&self) -> String {
1508        format!("{}:{}", self.avenue.kind, self.avenue.id)
1509    }
1510
1511    async fn request_unit(
1512        &self,
1513        build: impl FnOnce(oneshot::Sender<Result<()>>) -> Command,
1514    ) -> Result<()> {
1515        let (reply, response) = oneshot::channel();
1516        self.commands
1517            .send(build(reply))
1518            .await
1519            .map_err(|_| anyhow!("coordination gateway actor stopped"))?;
1520        response
1521            .await
1522            .map_err(|_| anyhow!("coordination gateway actor dropped its reply"))?
1523    }
1524
1525    async fn mint_credential(
1526        &self,
1527        desired: &DesiredPresence,
1528        purpose: &str,
1529        refresh_grant: Option<String>,
1530    ) -> Result<NativeGatewayGrant> {
1531        let ticket_fingerprint = base64::engine::general_purpose::URL_SAFE_NO_PAD
1532            .encode(Sha256::digest(desired.ticket.as_bytes()));
1533        let credential = self
1534            .grant_provider
1535            .grant(NativeGatewayGrantRequest {
1536                avenue: self.avenue.clone(),
1537                device_id: self.device_id.clone(),
1538                runtime_instance_id: self.runtime_instance_id.clone(),
1539                ticket_fingerprint,
1540                purpose: purpose.to_string(),
1541                architecture: self.architecture,
1542                room_delivery: self.room_delivery,
1543                refresh_grant,
1544            })
1545            .await
1546            .context("obtain native coordination grant")?;
1547        if credential.protocol_version != GRANT_PROTOCOL_VERSION {
1548            return Err(gateway_connect_error(
1549                "unsupported native coordination protocol",
1550                false,
1551            ));
1552        }
1553        let configured = reqwest::Url::parse(&self.endpoint)?;
1554        let returned = reqwest::Url::parse(&credential.gateway_url)?;
1555        if normalized_origin_scheme(configured.scheme())
1556            != normalized_origin_scheme(returned.scheme())
1557            || configured.host_str() != returned.host_str()
1558            || configured.port_or_known_default() != returned.port_or_known_default()
1559        {
1560            return Err(gateway_connect_error(
1561                "native coordination credential returned an unexpected gateway origin",
1562                false,
1563            ));
1564        }
1565        Ok(credential)
1566    }
1567
1568    fn gateway_url(&self, credential: &NativeGatewayGrant) -> Result<String> {
1569        let mut url = reqwest::Url::parse(&credential.gateway_url)?;
1570        match url.scheme() {
1571            "https" => url
1572                .set_scheme("wss")
1573                .map_err(|_| anyhow!("invalid gateway scheme"))?,
1574            "http" => url
1575                .set_scheme("ws")
1576                .map_err(|_| anyhow!("invalid gateway scheme"))?,
1577            "wss" | "ws" => {}
1578            _ => bail!("native coordination gateway must use HTTPS/WSS"),
1579        }
1580        if url.query().is_some() || url.fragment().is_some() || !url.username().is_empty() {
1581            bail!("native coordination gateway URL contains forbidden credentials or parameters");
1582        }
1583        let path = format!(
1584            "{}/v{}/connect/{}",
1585            url.path().trim_end_matches('/'),
1586            credential.protocol_version,
1587            credential.route_key
1588        );
1589        url.set_path(&path);
1590        Ok(url.to_string())
1591    }
1592
1593    fn device_from_gateway(&self, value: GatewayDevice) -> Device {
1594        Device {
1595            app_tag: Some(self.app_tag.clone()),
1596            device_id: value.device_id,
1597            user_id: value.user_id,
1598            device_name: value.device_name,
1599            platform_type: Some(value.platform_type),
1600            capabilities: value.capabilities,
1601            session_id: Some(value.runtime_instance_id),
1602            node_id: Some(value.node_id),
1603            tag: None,
1604            kind: None,
1605            metadata: value.metadata,
1606            online: value.online,
1607            ticket: Some(value.ticket),
1608            last_seen_at: Some(serde_json::json!(value.updated_at_ms)),
1609            expires_at: Some(serde_json::json!(value.expires_at_ms)),
1610            created_at: None,
1611            updated_at: Some(serde_json::json!(value.updated_at_ms)),
1612            excluded_peers: value.excluded_peers,
1613        }
1614    }
1615
1616    fn member_from_gateway(&self, value: GatewayMember) -> Device {
1617        Device {
1618            app_tag: Some(self.app_tag.clone()),
1619            device_id: value.device_id,
1620            user_id: value.user_id,
1621            device_name: value.device_name,
1622            platform_type: Some(value.platform_type),
1623            capabilities: value.capabilities,
1624            session_id: None,
1625            node_id: None,
1626            tag: None,
1627            kind: None,
1628            metadata: value.metadata,
1629            online: value.online,
1630            ticket: None,
1631            last_seen_at: Some(serde_json::json!(value.updated_at_ms)),
1632            expires_at: Some(serde_json::json!(value.expires_at_ms)),
1633            created_at: None,
1634            updated_at: Some(serde_json::json!(value.updated_at_ms)),
1635            excluded_peers: Vec::new(),
1636        }
1637    }
1638}
1639
1640async fn rebuild_sparse_devices(adapter: &NativeCoordinationGatewaySignaling) {
1641    let members = adapter.shared.members.read().await.clone();
1642    let routes = {
1643        let mut routes = adapter.shared.topology_routes.write().await;
1644        // A later membership re-add must await a new private route lease, not
1645        // resurrect a ticket retained before that member was withdrawn.
1646        routes.retain(|device_id, _| members.contains_key(device_id));
1647        routes.clone()
1648    };
1649    let mut next = HashMap::new();
1650    for member in members.into_values() {
1651        let device_id = member.device_id.clone();
1652        let device = routes
1653            .get(&device_id)
1654            .cloned()
1655            .map(|(route, _active)| adapter.device_from_gateway(route))
1656            .unwrap_or_else(|| adapter.member_from_gateway(member));
1657        next.insert(device_id, device);
1658    }
1659    let mut current = adapter.shared.devices.write().await;
1660    let mut events = current
1661        .keys()
1662        .filter(|device_id| !next.contains_key(*device_id))
1663        .cloned()
1664        .map(|device_id| DeviceEvent::Removed { device_id })
1665        .collect::<Vec<_>>();
1666    for (device_id, device) in &next {
1667        match current.get(device_id) {
1668            None => events.push(DeviceEvent::Added {
1669                device: device.clone(),
1670            }),
1671            Some(existing) if existing != device => events.push(DeviceEvent::Modified {
1672                device: device.clone(),
1673            }),
1674            Some(_) => {}
1675        }
1676    }
1677    *current = next;
1678    drop(current);
1679    if !events.is_empty() {
1680        let _ = adapter.shared.device_events.send(events);
1681    }
1682}
1683
1684type GatewaySocket =
1685    tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
1686
1687struct ConnectedGateway {
1688    socket: GatewaySocket,
1689    credential_expires_at_ms: u64,
1690    credential_token: String,
1691    lease_refresh_mode: LeaseRefreshMode,
1692}
1693
1694#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1695enum LeaseRefreshMode {
1696    Off,
1697    Shadow,
1698    Active,
1699}
1700
1701impl LeaseRefreshMode {
1702    fn from_wire(value: Option<&str>) -> Self {
1703        match value {
1704            Some("shadow") => Self::Shadow,
1705            Some("active") => Self::Active,
1706            _ => Self::Off,
1707        }
1708    }
1709}
1710
1711async fn run_actor(
1712    adapter: Arc<NativeCoordinationGatewaySignaling>,
1713    mut commands: mpsc::Receiver<Command>,
1714) {
1715    let mut socket: Option<ConnectedGateway> = None;
1716    let mut reconnect_attempt = 0_u8;
1717    let mut authority_return_retry_available = false;
1718    let mut retry_at: Option<tokio::time::Instant> = None;
1719    let mut pending_publish: Option<(DesiredPresence, oneshot::Sender<Result<()>>)> = None;
1720    let mut circuit_failure: Option<(DesiredPresence, anyhow::Error)> = None;
1721    let mut keepalive = tokio::time::interval(SOCKET_KEEPALIVE_INTERVAL);
1722    keepalive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1723    loop {
1724        if socket.is_none() {
1725            if pending_publish
1726                .as_ref()
1727                .is_some_and(|(_, reply)| reply.is_closed())
1728            {
1729                pending_publish = None;
1730                retry_at = None;
1731                reconnect_attempt = 0;
1732                *adapter.shared.desired.write().await = None;
1733                *adapter.shared.applied_presence.write().await = None;
1734            }
1735            let retry_delay = retry_at
1736                .map(|deadline| deadline.saturating_duration_since(tokio::time::Instant::now()))
1737                .unwrap_or(Duration::from_secs(365 * 24 * 60 * 60));
1738            tokio::select! {
1739                command = commands.recv() => {
1740                    let Some(command) = command else { break; };
1741                    match command {
1742                        Command::Publish { desired, reply } => {
1743                            if let Some((blocked, message)) = circuit_failure.as_ref() {
1744                                if blocked == &desired {
1745                                    let _ = reply.send(Err(copy_gateway_failure(message)));
1746                                    continue;
1747                                }
1748                            }
1749                            circuit_failure = None;
1750                            if let Some((_, previous_reply)) = pending_publish.take() {
1751                                let _ = previous_reply.send(Err(anyhow!(
1752                                    "native coordination publication was superseded"
1753                                )));
1754                            }
1755                            *adapter.shared.desired.write().await = Some(desired.clone());
1756                            *adapter.shared.applied_presence.write().await = None;
1757                            pending_publish = Some((desired, reply));
1758                            // New desired presence cannot bypass an in-flight
1759                            // admission cooldown or reset its attempt budget.
1760                            if retry_at.is_none() {
1761                                reconnect_attempt = 0;
1762                                retry_at = Some(tokio::time::Instant::now());
1763                            }
1764                        }
1765                        Command::Stop => {
1766                            if let Some((_, reply)) = pending_publish.take() {
1767                                let _ = reply.send(Err(anyhow!(
1768                                    "native coordination gateway stopped"
1769                                )));
1770                            }
1771                            break;
1772                        }
1773                        Command::Delete { user_id, device_id, reply } => {
1774                            let result = execute_control_delete(
1775                                &adapter,
1776                                &user_id,
1777                                &device_id,
1778                            )
1779                            .await;
1780                            let _ = reply.send(result);
1781                        }
1782                        other => reject_command(
1783                            other,
1784                            if retry_at.is_some() {
1785                                "native coordination gateway is reconnecting"
1786                            } else {
1787                                "native coordination gateway is not connected"
1788                            },
1789                        ),
1790                    }
1791                }
1792                _ = tokio::time::sleep(retry_delay), if retry_at.is_some() => {
1793                    retry_at = None;
1794                    let desired = adapter.shared.desired.read().await.clone();
1795                    let Some(desired) = desired else {
1796                        reconnect_attempt = 0;
1797                        continue;
1798                    };
1799                    // Gateway authentication validates and persists the full
1800                    // device projection. `ready` is therefore the initial
1801                    // publication acknowledgement; sending presence.upsert
1802                    // here would duplicate the write and customer charge.
1803                    let result = connect(&adapter, pending_publish.as_mut().map(|(desired, _)| desired)).await;
1804                    match result {
1805                        Ok(connected) => {
1806                            reconnect_attempt = 0;
1807                            authority_return_retry_available = true;
1808                            circuit_failure = None;
1809                            // Connect may renew an expired locally owned room
1810                            // ticket before minting its fingerprint-bound grant.
1811                            *adapter.shared.applied_presence.write().await = adapter.shared.desired.read().await.clone();
1812                            socket = Some(connected);
1813                            if let Some((published, reply)) = pending_publish.take() {
1814                                if Some(&published) == adapter.shared.applied_presence.read().await.as_ref() {
1815                                    let _ = reply.send(Ok(()));
1816                                } else {
1817                                    let _ = reply.send(Err(anyhow!(
1818                                        "native coordination publication was superseded"
1819                                    )));
1820                                }
1821                            }
1822                        }
1823                        Err(error) => {
1824                            reconnect_attempt = reconnect_attempt.saturating_add(1);
1825                            let retryable = adapter.handle_gateway_failure(&error).await;
1826                            eprintln!(
1827                                "[openrtc][coordination-gateway][connect-retry] attempt={}/{} retryable={} error={}",
1828                                reconnect_attempt,
1829                                MAX_RECONNECT_ATTEMPTS,
1830                                retryable,
1831                                error
1832                            );
1833                            let delay = gateway_reconnect_delay(
1834                                reconnect_attempt, &error, &mut authority_return_retry_available,
1835                            );
1836                            if let Some(delay) = delay {
1837                                retry_at = Some(tokio::time::Instant::now() + delay);
1838                            } else {
1839                                retry_at = None;
1840                                let failure = error.context(format!(
1841                                    "native coordination gateway unavailable after {} attempts",
1842                                    reconnect_attempt
1843                                ));
1844                                if let Some((_, reply)) = pending_publish.take() {
1845                                    let _ = reply.send(Err(copy_gateway_failure(&failure)));
1846                                }
1847                                circuit_failure = Some((desired, failure));
1848                            }
1849                        }
1850                    }
1851                }
1852                _ = tokio::time::sleep(adapter.authority_ticket_maintenance_delay().await) => {
1853                    // Ticket lifetime is independent of gateway retry/backoff.
1854                    // This is local maintenance in the same actor: it cannot
1855                    // mint a gateway grant or reset the network attempt budget.
1856                    if let Err(error) = adapter.maintain_authority_publication(
1857                        pending_publish.as_mut().map(|(desired, _)| desired),
1858                    ).await {
1859                        eprintln!("[openrtc][coordination-gateway][offline-ticket-failed] error={error:#}");
1860                        adapter.stop_authority_client().await;
1861                        retry_at = None;
1862                        if let Some(desired) = adapter.shared.desired.read().await.clone() {
1863                            circuit_failure = Some((desired, copy_gateway_failure(&error)));
1864                        }
1865                        if let Some((_, reply)) = pending_publish.take() {
1866                            let _ = reply.send(Err(error));
1867                        }
1868                    }
1869                }
1870            }
1871            continue;
1872        }
1873
1874        if let Some(active) = socket.as_mut() {
1875            tokio::select! {
1876                command = commands.recv() => {
1877                    let Some(command) = command else { break; };
1878                    if matches!(command, Command::Stop) {
1879                        let _ = active.socket.close(None).await;
1880                        break;
1881                    }
1882                    if let Some((blocked, message)) = circuit_failure.as_ref() {
1883                        if matches!(
1884                            &command,
1885                            Command::Publish { desired, .. } if desired == blocked
1886                        ) {
1887                            let Command::Publish { reply, .. } = command else {
1888                                unreachable!("only unchanged publication enters this branch")
1889                            };
1890                            let _ = reply.send(Err(copy_gateway_failure(message)));
1891                            continue;
1892                        }
1893                    }
1894                    let credential_rebind = match &command {
1895                        Command::Publish { desired, .. } => adapter
1896                            .shared
1897                            .applied_presence
1898                            .read()
1899                            .await
1900                            .as_ref()
1901                            .is_some_and(|applied| applied.ticket != desired.ticket),
1902                        _ => false,
1903                    };
1904                    if credential_rebind {
1905                        let Command::Publish { desired, reply } = command else {
1906                            unreachable!("credential rebind is only set for publication")
1907                        };
1908                        *adapter.shared.desired.write().await = Some(desired.clone());
1909                        *adapter.shared.applied_presence.write().await = None;
1910                        pending_publish = Some((desired, reply));
1911                        let _ = active.socket.close(None).await;
1912                        socket = None;
1913                        reconnect_attempt = 0;
1914                        retry_at = Some(tokio::time::Instant::now());
1915                        continue;
1916                    }
1917                    let going_offline = matches!(&command, Command::Offline { .. });
1918                    let publishing = matches!(&command, Command::Publish { .. });
1919                    let deleting_local_device =
1920                        command_deletes_device(&command, &adapter.device_id);
1921                    let result = handle_command(&adapter, &mut active.socket, command).await;
1922                    if going_offline || (deleting_local_device && result.is_ok()) {
1923                        let _ = active.socket.close(None).await;
1924                        *adapter.shared.desired.write().await = None;
1925                        *adapter.shared.applied_presence.write().await = None;
1926                        socket = None;
1927                        reconnect_attempt = 0;
1928                        retry_at = None;
1929                        circuit_failure = None;
1930                    } else if let Err(error) = result {
1931                        if adapter.handle_gateway_failure(&error).await {
1932                            *adapter.shared.applied_presence.write().await = None;
1933                            socket = None;
1934                            reconnect_attempt = 0;
1935                            retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
1936                        } else {
1937                            eprintln!(
1938                                "[openrtc][coordination-gateway][operation-rejected] retryable=false error={}",
1939                                error
1940                            );
1941                            if is_gateway_revocation(&error) {
1942                                let _ = active.socket.close(None).await;
1943                                *adapter.shared.applied_presence.write().await = None;
1944                                socket = None;
1945                                retry_at = None;
1946                                if let Some(desired) = adapter.shared.desired.read().await.clone() {
1947                                    circuit_failure = Some((desired, error));
1948                                }
1949                                continue;
1950                            }
1951                            if publishing {
1952                                if let Some(desired) = adapter.shared.desired.read().await.clone() {
1953                                    circuit_failure = Some((desired, error));
1954                                }
1955                            }
1956                        }
1957                    }
1958                }
1959                _ = tokio::time::sleep(adapter.authority_ticket_maintenance_delay().await) => {
1960                    let renewal = async {
1961                        if adapter.maintain_authority_ticket().await? {
1962                            let desired = adapter.shared.desired.read().await.clone()
1963                                .ok_or_else(|| anyhow!("authority presence was withdrawn"))?;
1964                            refresh_presence_authentication(&adapter, active, desired).await?;
1965                        }
1966                        Ok::<_, anyhow::Error>(())
1967                    }.await;
1968                    if let Err(error) = renewal {
1969                        let retryable = adapter.handle_gateway_failure(&error).await;
1970                        eprintln!("[openrtc][coordination-gateway][authority-ticket-failed] retryable={} error={error:#}", retryable);
1971                        *adapter.shared.applied_presence.write().await = None;
1972                        let _ = active.socket.close(None).await;
1973                        socket = None;
1974                        reconnect_attempt = 0;
1975                        if retryable {
1976                            retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
1977                        } else {
1978                            retry_at = None;
1979                            if let Some(desired) = adapter.shared.desired.read().await.clone() {
1980                                circuit_failure = Some((desired, error));
1981                            }
1982                        }
1983                    }
1984                }
1985                _ = tokio::time::sleep(auth_refresh_delay(active.credential_expires_at_ms)) => {
1986                    let refresh = if active.lease_refresh_mode == LeaseRefreshMode::Active {
1987                        refresh_lease(&adapter, &mut active.socket).await
1988                    } else {
1989                        refresh_authentication(
1990                            &adapter,
1991                            &mut active.socket,
1992                            Some(active.credential_token.clone()),
1993                        ).await
1994                    };
1995                    match refresh {
1996                        Ok((expires_at_ms, credential_token)) => {
1997                            active.credential_expires_at_ms = expires_at_ms;
1998                            if let Some(credential_token) = credential_token {
1999                                active.credential_token = credential_token;
2000                            }
2001                        }
2002                        Err(error) => {
2003                            let terminal_error = if active.lease_refresh_mode == LeaseRefreshMode::Active
2004                                && !is_gateway_revocation(&error)
2005                                && gateway_service_error(&error).is_none() {
2006                                eprintln!(
2007                                    "[openrtc][coordination-gateway][lease-refresh-failed] retryable={} error={}",
2008                                    is_retryable_gateway_error(&error),
2009                                    error,
2010                                );
2011                                match refresh_authentication(
2012                                    &adapter,
2013                                    &mut active.socket,
2014                                    None,
2015                                ).await {
2016                                    Ok((expires_at_ms, Some(token))) => {
2017                                        active.credential_expires_at_ms = expires_at_ms;
2018                                        active.credential_token = token;
2019                                        continue;
2020                                    }
2021                                    Ok(_) => terminal_refresh_error(error, None),
2022                                    Err(fallback) => {
2023                                        eprintln!(
2024                                            "[openrtc][coordination-gateway][lease-fallback-failed] retryable={} error={}",
2025                                            is_retryable_gateway_error(&fallback),
2026                                            fallback,
2027                                        );
2028                                        terminal_refresh_error(error, Some(fallback))
2029                                    }
2030                                }
2031                            } else {
2032                                error
2033                            };
2034                            let retryable = adapter.handle_gateway_failure(&terminal_error).await;
2035                            eprintln!(
2036                                "[openrtc][coordination-gateway][auth-refresh-failed] retryable={} error={}",
2037                                retryable, terminal_error,
2038                            );
2039                            *adapter.shared.applied_presence.write().await = None;
2040                            socket = None;
2041                            reconnect_attempt = 0;
2042                            if retryable {
2043                                retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&terminal_error));
2044                            } else {
2045                                retry_at = None;
2046                                if let Some(desired) = adapter.shared.desired.read().await.clone() {
2047                                    circuit_failure = Some((desired, terminal_error));
2048                                }
2049                            }
2050                        }
2051                    }
2052                }
2053                _ = keepalive.tick() => {
2054                    if let Err(error) = send_socket_keepalive(&mut active.socket).await {
2055                        eprintln!(
2056                            "[openrtc][coordination-gateway][keepalive-failed] retryable=true error={}",
2057                            error,
2058                        );
2059                        *adapter.shared.applied_presence.write().await = None;
2060                        socket = None;
2061                        reconnect_attempt = 0;
2062                        retry_at = Some(tokio::time::Instant::now());
2063                    }
2064                }
2065                incoming = active.socket.next() => {
2066                    match incoming {
2067                        Some(Ok(message)) => {
2068                            if let Err(error) = handle_message(&adapter, &mut active.socket, message).await {
2069                                let retryable = adapter.handle_gateway_failure(&error).await;
2070                                eprintln!(
2071                                    "[openrtc][coordination-gateway][socket-frame-failed] retryable={} error={:#}",
2072                                    retryable,
2073                                    error,
2074                                );
2075                                *adapter.shared.applied_presence.write().await = None;
2076                                socket = None;
2077                                reconnect_attempt = 0;
2078                                if retryable {
2079                                    retry_at = Some(tokio::time::Instant::now() + gateway_service_retry_delay(&error));
2080                                } else {
2081                                    retry_at = None;
2082                                    if let Some(desired) = adapter.shared.desired.read().await.clone() {
2083                                        circuit_failure = Some((desired, error));
2084                                    }
2085                                }
2086                            }
2087                        }
2088                        Some(Err(error)) => {
2089                            eprintln!(
2090                                "[openrtc][coordination-gateway][socket-receive-failed] retryable=true error={:#}",
2091                                error,
2092                            );
2093                            *adapter.shared.applied_presence.write().await = None;
2094                            socket = None;
2095                            reconnect_attempt = 0;
2096                            retry_at = Some(tokio::time::Instant::now());
2097                        }
2098                        None => {
2099                            eprintln!(
2100                                "[openrtc][coordination-gateway][socket-closed] retryable=true"
2101                            );
2102                            *adapter.shared.applied_presence.write().await = None;
2103                            socket = None;
2104                            reconnect_attempt = 0;
2105                            retry_at = Some(tokio::time::Instant::now());
2106                        },
2107                    }
2108                }
2109            }
2110        }
2111    }
2112    adapter.stop_authority_client().await;
2113}
2114
2115fn command_deletes_device(command: &Command, local_device_id: &str) -> bool {
2116    matches!(
2117        command,
2118        Command::Delete { device_id, .. } if device_id == local_device_id
2119    )
2120}
2121
2122async fn connect(adapter: &NativeCoordinationGatewaySignaling, pending: Option<&mut DesiredPresence>) -> Result<ConnectedGateway> {
2123    adapter.maintain_authority_publication(pending).await?;
2124    let desired = adapter
2125        .shared
2126        .desired
2127        .read()
2128        .await
2129        .clone()
2130        .ok_or_else(|| anyhow!("native coordination presence is not configured"))?;
2131    let work = connect_with_desired(adapter, &desired, "presence");
2132    tokio::pin!(work);
2133    let mut expiry_delay = adapter.expire_retiring_authority_tickets().await;
2134    loop {
2135        tokio::select! {
2136            result = &mut work => return result,
2137            delay = async {
2138                tokio::time::sleep(expiry_delay).await;
2139                adapter.expire_retiring_authority_tickets().await
2140            } => expiry_delay = delay,
2141        }
2142    }
2143}
2144
2145async fn connect_with_desired(
2146    adapter: &NativeCoordinationGatewaySignaling,
2147    desired: &DesiredPresence,
2148    purpose: &str,
2149) -> Result<ConnectedGateway> {
2150    let credential = adapter.mint_credential(desired, purpose, None).await?;
2151    if credential.expires_at_ms <= now_ms().saturating_add(30_000) {
2152        bail!("native coordination credential expires too soon");
2153    }
2154    if purpose == "presence" {
2155        *adapter.shared.topology_revision.write().await = 0;
2156        *adapter.shared.active_grant.write().await = Some((
2157            gateway_grant_jti(&credential.token)?,
2158            credential.expires_at_ms,
2159        ));
2160    }
2161    let gateway_url = adapter.gateway_url(&credential)?;
2162    let uri: Uri = gateway_url.parse().context("parse native gateway URI")?;
2163    let mut request = uri.into_client_request()?;
2164    request.headers_mut().insert(
2165        "Sec-WebSocket-Protocol",
2166        HeaderValue::from_str(&format!(
2167            "{}, {}{}",
2168            GATEWAY_PROTOCOL, GATEWAY_AUTH_PROTOCOL_PREFIX, credential.token
2169        ))?,
2170    );
2171    let (mut socket, _) =
2172        tokio::time::timeout(REQUEST_TIMEOUT, tokio_tungstenite::connect_async(request))
2173            .await
2174            .context("native coordination WebSocket connect timed out")??;
2175    socket
2176        .send(Message::Text(
2177            serde_json::json!({
2178                "v": WIRE_PROTOCOL_VERSION,
2179                "type": "auth",
2180                "token": credential.token,
2181                "socketLiveness": "ping-v1",
2182                "device": gateway_device_value(adapter, &desired).await
2183            })
2184            .to_string()
2185            .into(),
2186        ))
2187        .await?;
2188    let lease_refresh_mode = tokio::time::timeout(ACK_TIMEOUT, async {
2189        while let Some(message) = socket.next().await {
2190            let message = message?;
2191            if let Some(frame) = decode_frame(message)? {
2192                match frame {
2193                    ServerFrame::Ready {
2194                        lease_refresh_mode, ..
2195                    } => {
2196                        return Ok::<LeaseRefreshMode, anyhow::Error>(LeaseRefreshMode::from_wire(
2197                            lease_refresh_mode.as_deref(),
2198                        ))
2199                    }
2200                    ServerFrame::Error {
2201                        code,
2202                        message,
2203                        retryable,
2204                        metadata,
2205                        ..
2206                    } => {
2207                        return Err(gateway_server_error_metadata(code, message, retryable, metadata));
2208                    }
2209                    other => handle_frame_with_socket(adapter, &mut socket, other).await?,
2210                }
2211            }
2212        }
2213        bail!("coordination gateway closed before ready")
2214    })
2215    .await
2216    .context("native coordination authentication timed out")??;
2217    Ok(ConnectedGateway {
2218        socket,
2219        credential_expires_at_ms: credential.expires_at_ms,
2220        credential_token: credential.token,
2221        lease_refresh_mode,
2222    })
2223}
2224
2225async fn execute_control_delete(
2226    adapter: &NativeCoordinationGatewaySignaling,
2227    user_id: &str,
2228    target_device_id: &str,
2229) -> Result<()> {
2230    let desired = DesiredPresence {
2231        user_id: user_id.to_string(),
2232        local_node_id: adapter.device_id.clone(),
2233        ticket: format!("openrtc-device-control:{}", adapter.runtime_instance_id),
2234        device_name: adapter.device_id.clone(),
2235        metadata: None,
2236        ttl_ms: 0,
2237        online: false,
2238    };
2239    let mut connected = connect_with_desired(adapter, &desired, "device-control").await?;
2240    let result = execute_operation(
2241        adapter,
2242        &mut connected.socket,
2243        serde_json::json!({
2244            "v": WIRE_PROTOCOL_VERSION,
2245            "type": "device.delete",
2246            "idempotencyKey": random_id("device-delete"),
2247            "targetDeviceId": target_device_id,
2248        }),
2249    )
2250    .await;
2251    let _ = connected.socket.close(None).await;
2252    result
2253}
2254
2255async fn refresh_authentication(
2256    adapter: &NativeCoordinationGatewaySignaling,
2257    socket: &mut GatewaySocket,
2258    current_grant: Option<String>,
2259) -> Result<(u64, Option<String>)> {
2260    let desired = adapter
2261        .shared
2262        .desired
2263        .read()
2264        .await
2265        .clone()
2266        .ok_or_else(|| anyhow!("native coordination presence is not configured"))?;
2267    refresh_authentication_for_presence(adapter, socket, current_grant, &desired).await
2268}
2269
2270async fn refresh_authentication_for_presence(
2271    adapter: &NativeCoordinationGatewaySignaling,
2272    socket: &mut GatewaySocket,
2273    current_grant: Option<String>,
2274    desired: &DesiredPresence,
2275) -> Result<(u64, Option<String>)> {
2276    let credential = adapter
2277        .mint_credential(desired, "presence", current_grant)
2278        .await?;
2279    let refreshed_token = credential.token.clone();
2280    let refreshed_jti = gateway_grant_jti(&credential.token)?;
2281    socket
2282        .send(Message::Text(
2283            serde_json::json!({
2284                "v": WIRE_PROTOCOL_VERSION,
2285                "type": "auth.refresh",
2286                "token": credential.token,
2287            })
2288            .to_string()
2289            .into(),
2290        ))
2291        .await?;
2292    tokio::time::timeout(ACK_TIMEOUT, async {
2293        while let Some(message) = socket.next().await {
2294            let message = message?;
2295            if let Some(frame) = decode_frame(message)? {
2296                match frame {
2297                    ServerFrame::AuthRefreshed { _expires_at_ms } => {
2298                        *adapter.shared.active_grant.write().await =
2299                            Some((refreshed_jti, _expires_at_ms));
2300                        update_local_device_expiry(adapter, _expires_at_ms).await;
2301                        return Ok((_expires_at_ms, Some(refreshed_token)));
2302                    }
2303                    ServerFrame::Error {
2304                        code,
2305                        message,
2306                        retryable,
2307                        metadata,
2308                        ..
2309                    } => {
2310                        return Err(gateway_server_error_metadata(code, message, retryable, metadata));
2311                    }
2312                    other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
2313                }
2314            }
2315        }
2316        bail!("coordination gateway closed before authentication refresh")
2317    })
2318    .await
2319    .context("native coordination authentication refresh timed out")?
2320}
2321
2322/// Rebind an expiring service ticket on the current gateway socket. Grant and
2323/// presence must describe the same snapshot, even if a caller stages a newer
2324/// publication while authorization is pending. Neither ACK implies peer trust.
2325async fn refresh_presence_authentication(
2326    adapter: &NativeCoordinationGatewaySignaling,
2327    active: &mut ConnectedGateway,
2328    desired: DesiredPresence,
2329) -> Result<()> {
2330    let (expires_at_ms, token) = refresh_authentication_for_presence(
2331        adapter,
2332        &mut active.socket,
2333        Some(active.credential_token.clone()),
2334        &desired,
2335    )
2336    .await?;
2337    active.credential_expires_at_ms = expires_at_ms;
2338    if let Some(token) = token {
2339        active.credential_token = token;
2340    }
2341    execute_operation_for_presence(
2342        adapter,
2343        &mut active.socket,
2344        presence_frame(adapter, &desired).await,
2345        Some(&desired),
2346    )
2347    .await?;
2348    *adapter.shared.applied_presence.write().await = Some(desired);
2349    Ok(())
2350}
2351
2352async fn refresh_lease(
2353    adapter: &NativeCoordinationGatewaySignaling,
2354    socket: &mut GatewaySocket,
2355) -> Result<(u64, Option<String>)> {
2356    let idempotency_key = random_id("lease");
2357    socket
2358        .send(Message::Text(
2359            serde_json::json!({
2360                "v": WIRE_PROTOCOL_VERSION,
2361                "type": "lease.refresh",
2362                "idempotencyKey": idempotency_key,
2363            })
2364            .to_string()
2365            .into(),
2366        ))
2367        .await?;
2368    tokio::time::timeout(ACK_TIMEOUT, async {
2369        while let Some(message) = socket.next().await {
2370            let message = message?;
2371            if let Some(frame) = decode_frame(message)? {
2372                match frame {
2373                    ServerFrame::LeaseRefreshed {
2374                        idempotency_key: response_key,
2375                        expires_at_ms,
2376                    } if response_key == idempotency_key => {
2377                        let mut grant = adapter.shared.active_grant.write().await;
2378                        let (_, current_expiry) = grant.as_mut().ok_or_else(|| {
2379                            gateway_connect_error("lease refresh has no current grant", false)
2380                        })?;
2381                        if expires_at_ms <= now_ms() {
2382                            return Err(gateway_connect_error(
2383                                "lease refresh is already expired",
2384                                false,
2385                            ));
2386                        }
2387                        // This ACK renews only the authenticated socket authority.
2388                        // Receiver permissions still require a fresh topology lease.
2389                        *current_expiry = expires_at_ms;
2390                        drop(grant);
2391                        update_local_device_expiry(adapter, expires_at_ms).await;
2392                        return Ok((expires_at_ms, None));
2393                    }
2394                    ServerFrame::Error {
2395                        code,
2396                        message,
2397                        retryable,
2398                        metadata,
2399                        ..
2400                    } => {
2401                        return Err(gateway_server_error_metadata(code, message, retryable, metadata));
2402                    }
2403                    other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
2404                }
2405            }
2406        }
2407        bail!("coordination gateway closed before lease refresh acknowledgement")
2408    })
2409    .await
2410    .context("native coordination lease refresh timed out")?
2411}
2412
2413async fn update_local_device_expiry(
2414    adapter: &NativeCoordinationGatewaySignaling,
2415    expires_at_ms: u64,
2416) {
2417    let event = {
2418        let mut devices = adapter.shared.devices.write().await;
2419        let Some(device) = devices.get_mut(&adapter.device_id) else {
2420            return;
2421        };
2422        let next_expiry = serde_json::json!(expires_at_ms);
2423        if device.expires_at.as_ref() == Some(&next_expiry) {
2424            return;
2425        }
2426        device.expires_at = Some(next_expiry);
2427        DeviceEvent::Modified {
2428            device: device.clone(),
2429        }
2430    };
2431    let _ = adapter.shared.device_events.send(vec![event]);
2432}
2433
2434async fn handle_command(
2435    adapter: &NativeCoordinationGatewaySignaling,
2436    socket: &mut GatewaySocket,
2437    command: Command,
2438) -> Result<()> {
2439    match command {
2440        Command::Publish { desired, reply } => {
2441            *adapter.shared.desired.write().await = Some(desired.clone());
2442            if adapter.shared.applied_presence.read().await.as_ref() == Some(&desired) {
2443                let _ = reply.send(Ok(()));
2444                return Ok(());
2445            }
2446            let result =
2447                execute_operation(adapter, socket, presence_frame(adapter, &desired).await).await;
2448            if result.is_ok() {
2449                *adapter.shared.applied_presence.write().await = Some(desired);
2450            }
2451            complete_command(reply, result, "native presence publication")?;
2452        }
2453        Command::Patch { patch, reply } => {
2454            merge_patch(
2455                &mut *adapter.shared.staged_patch.write().await,
2456                patch.clone(),
2457            );
2458            let result = execute_operation(
2459                adapter,
2460                socket,
2461                serde_json::json!({
2462                    "v": WIRE_PROTOCOL_VERSION,
2463                    "type": "device.patch",
2464                    "idempotencyKey": random_id("device-patch"),
2465                    "patch": patch,
2466                }),
2467            )
2468            .await;
2469            complete_command(reply, result, "native device patch")?;
2470        }
2471        Command::Offline { reply } => {
2472            if let Some(desired) = adapter.shared.desired.write().await.as_mut() {
2473                desired.online = false;
2474            }
2475            let result = execute_operation(
2476                adapter,
2477                socket,
2478                serde_json::json!({
2479                    "v": WIRE_PROTOCOL_VERSION,
2480                    "type": "presence.offline",
2481                    "idempotencyKey": random_id("offline"),
2482                }),
2483            )
2484            .await;
2485            complete_command(reply, result, "native presence offline")?;
2486        }
2487        Command::Delete {
2488            device_id, reply, ..
2489        } => {
2490            let result = execute_operation(
2491                adapter,
2492                socket,
2493                serde_json::json!({
2494                    "v": WIRE_PROTOCOL_VERSION,
2495                    "type": "device.delete",
2496                    "idempotencyKey": random_id("device-delete"),
2497                    "targetDeviceId": device_id,
2498                }),
2499            )
2500            .await;
2501            complete_command(reply, result, "native device delete")?;
2502        }
2503        Command::SendSignal {
2504            target_device_id,
2505            payload,
2506            state,
2507            reply_payload,
2508            reply,
2509        } => {
2510            let idempotency_key = random_id("signal");
2511            let result = execute_operation(
2512                adapter,
2513                socket,
2514                serde_json::json!({
2515                    "v": WIRE_PROTOCOL_VERSION,
2516                    "type": "signal.send",
2517                    "idempotencyKey": idempotency_key,
2518                    "targetDeviceId": target_device_id,
2519                    "payload": payload,
2520                    "state": state,
2521                    "replyPayload": reply_payload,
2522                }),
2523            )
2524            .await
2525            .map(|_| idempotency_key);
2526            complete_command(reply, result, "native signal send")?;
2527        }
2528        Command::PutSession {
2529            session_id,
2530            session,
2531            expires_at_ms,
2532            reply,
2533        } => {
2534            let result = execute_operation(
2535                adapter,
2536                socket,
2537                serde_json::json!({
2538                    "v": WIRE_PROTOCOL_VERSION,
2539                    "type": "session.put",
2540                    "idempotencyKey": random_id("session-put"),
2541                    "sessionId": session_id,
2542                    "session": session,
2543                    "expiresAtMs": expires_at_ms,
2544                }),
2545            )
2546            .await;
2547            complete_command(reply, result, "native session update")?;
2548        }
2549        Command::PublishAuthorityAssignment {
2550            assignment,
2551            signature,
2552            reply,
2553        } => {
2554            let result = execute_operation(
2555                adapter,
2556                socket,
2557                serde_json::json!({
2558                    "v": WIRE_PROTOCOL_VERSION,
2559                    "type": "authority.assignment.publish",
2560                    "idempotencyKey": random_id("authority-assignment"),
2561                    "assignment": assignment,
2562                    "signature": signature,
2563                }),
2564            )
2565            .await;
2566            complete_command(reply, result, "native authority assignment publish")?;
2567        }
2568        #[cfg(feature = "managed-group-encryption")]
2569        Command::PublishManaged { batch, reply } => {
2570            let result = publish_native_managed_batch(adapter, socket, batch).await;
2571            complete_command(reply, result, "native managed room publish")?;
2572        }
2573        Command::Stop => {}
2574    }
2575    Ok(())
2576}
2577
2578fn complete_command<T>(
2579    reply: oneshot::Sender<Result<T>>,
2580    result: Result<T>,
2581    operation: &str,
2582) -> Result<()> {
2583    match result {
2584        Ok(value) => {
2585            let _ = reply.send(Ok(value));
2586            Ok(())
2587        }
2588        Err(error) => {
2589            // Preserve the owning runtime's full failure chain for the caller.
2590            // `anyhow::Error::to_string()` reports only the outer context, which
2591            // previously hid whether a failed budget renewal came from the
2592            // control plane, the gateway refresh frame, or the socket itself.
2593            let message = format!("{error:#}");
2594            let _ = reply.send(Err(copy_gateway_failure(&error)));
2595            Err(error.context(format!("{operation} failed: {message}")))
2596        }
2597    }
2598}
2599
2600async fn execute_operation(
2601    adapter: &NativeCoordinationGatewaySignaling,
2602    socket: &mut GatewaySocket,
2603    frame: serde_json::Value,
2604) -> Result<()> {
2605    execute_operation_for_presence(adapter, socket, frame, None).await
2606}
2607
2608async fn execute_operation_for_presence(
2609    adapter: &NativeCoordinationGatewaySignaling,
2610    socket: &mut GatewaySocket,
2611    frame: serde_json::Value,
2612    bound_presence: Option<&DesiredPresence>,
2613) -> Result<()> {
2614    let operation = frame
2615        .get("type")
2616        .and_then(serde_json::Value::as_str)
2617        .unwrap_or("unknown")
2618        .to_string();
2619    let idempotency_key = frame
2620        .get("idempotencyKey")
2621        .and_then(serde_json::Value::as_str)
2622        .ok_or_else(|| anyhow!("gateway operation is missing idempotencyKey"))?
2623        .to_string();
2624    for attempt in 0..MAX_OPERATION_ATTEMPTS {
2625        socket
2626            .send(Message::Text(frame.to_string().into()))
2627            .await
2628            .context("send native gateway operation")?;
2629        let renewal_required = tokio::time::timeout(ACK_TIMEOUT, async {
2630            while let Some(message) = socket.next().await {
2631                let message = message?;
2632                if let Some(frame) = decode_frame(message)? {
2633                    match frame {
2634                        ServerFrame::Ack {
2635                            idempotency_key: ack,
2636                        } if ack == idempotency_key => return Ok(false),
2637                        ServerFrame::Error {
2638                            code,
2639                            message,
2640                            idempotency_key: Some(key),
2641                            retryable,
2642                            metadata,
2643                        } if key == idempotency_key => {
2644                            if operation_requires_budget_renewal(&code, retryable) {
2645                                return Ok(true);
2646                            }
2647                            return Err(gateway_server_error_metadata(
2648                                code,
2649                                format!("{message} operation={operation}"),
2650                                retryable,
2651                                metadata,
2652                            ));
2653                        }
2654                        other => Box::pin(handle_frame_with_socket(adapter, socket, other)).await?,
2655                    }
2656                }
2657            }
2658            bail!("coordination gateway closed before acknowledgement")
2659        })
2660        .await
2661        .context("native coordination acknowledgement timed out")??;
2662        if !renewal_required {
2663            return Ok(());
2664        }
2665        if attempt > 0 {
2666            bail!("native coordination budget renewal did not fund operation={operation}");
2667        }
2668        // The browser SDK uses the same one-renewal/one-replay contract. Keep
2669        // the original idempotency key so an ambiguous first attempt cannot
2670        // create a second logical charge. A full grant refresh asks the
2671        // control plane for current budget truth instead of trusting the old
2672        // signed spend snapshot.
2673        let renewal = match bound_presence {
2674            Some(desired) => {
2675                refresh_authentication_for_presence(adapter, socket, None, desired).await
2676            }
2677            None => refresh_authentication(adapter, socket, None).await,
2678        };
2679        renewal
2680            .with_context(|| format!("renew native coordination budget operation={operation}"))?;
2681    }
2682    unreachable!("native coordination operation loop returns or errors")
2683}
2684
2685async fn handle_message(
2686    adapter: &NativeCoordinationGatewaySignaling,
2687    socket: &mut GatewaySocket,
2688    message: Message,
2689) -> Result<()> {
2690    if let Some(frame) = decode_frame(message)? {
2691        handle_frame_with_socket(adapter, socket, frame).await?;
2692    }
2693    Ok(())
2694}
2695
2696#[cfg(feature = "managed-group-encryption")]
2697async fn ensure_managed_group(
2698    adapter: &NativeCoordinationGatewaySignaling,
2699) -> Result<Option<ManagedGroupAction>> {
2700    if adapter.shared.managed_group.lock().await.is_some() {
2701        return Ok(None);
2702    }
2703    let signer = adapter
2704        .managed_group_signer
2705        .as_ref()
2706        .ok_or_else(|| anyhow!("native managed rooms require a sign-only device key"))?;
2707    let store = adapter
2708        .managed_group_store
2709        .as_ref()
2710        .ok_or_else(|| anyhow!("native managed rooms require durable app-private state storage"))?;
2711    let avenue_key = adapter.avenue_key();
2712    let challenge = format!(
2713        "openrtc:managed-group-wrapping-key:v1:{}:{}:{}",
2714        adapter.app_tag, adapter.device_id, avenue_key,
2715    );
2716    let mut signature = signer
2717        .sign(&adapter.app_tag, challenge.as_bytes())
2718        .context("derive native managed group wrapping key")?;
2719    if signature.len() < 32 {
2720        bail!("native managed group device signature is invalid");
2721    }
2722    let mut digest = Sha256::new();
2723    digest.update(b"openrtc:managed-group-key-derivation:v1:");
2724    digest.update(challenge.as_bytes());
2725    digest.update(&signature);
2726    let wrapping_key: [u8; 32] = digest.finalize().into();
2727    signature.zeroize();
2728    let sealed_state = store
2729        .load(&adapter.app_tag, &adapter.device_id, &avenue_key)
2730        .await
2731        .context("load native managed group state")?;
2732    let controller = ManagedGroupController::new(
2733        adapter.device_id.clone(),
2734        wrapping_key,
2735        sealed_state.as_deref(),
2736    )?;
2737    let action = controller.publish_key_package()?;
2738    let mut current = adapter.shared.managed_group.lock().await;
2739    if current.is_some() {
2740        return Ok(None);
2741    }
2742    *current = Some(controller);
2743    Ok(Some(action))
2744}
2745
2746#[cfg(feature = "managed-group-encryption")]
2747async fn execute_managed_actions(
2748    adapter: &NativeCoordinationGatewaySignaling,
2749    socket: &mut GatewaySocket,
2750    actions: Vec<ManagedGroupAction>,
2751) -> Result<()> {
2752    for action in actions {
2753        match action {
2754            ManagedGroupAction::PersistState { sealed_state } => {
2755                let state = base64::engine::general_purpose::URL_SAFE_NO_PAD
2756                    .decode(sealed_state)
2757                    .context("decode native managed group persisted state")?;
2758                adapter
2759                    .managed_group_store
2760                    .as_ref()
2761                    .ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
2762                    .save(
2763                        &adapter.app_tag,
2764                        &adapter.device_id,
2765                        &adapter.avenue_key(),
2766                        &state,
2767                    )
2768                    .await
2769                    .context("persist native managed group state")?;
2770            }
2771            ManagedGroupAction::PublishKeyPackage {
2772                key_package_hash,
2773                key_package,
2774            } => {
2775                execute_operation(
2776                    adapter,
2777                    socket,
2778                    serde_json::json!({
2779                        "v": WIRE_PROTOCOL_VERSION,
2780                        "type": "managed.key-package.publish",
2781                        "idempotencyKey": random_id("managed-key-package"),
2782                        "keyPackageHash": key_package_hash,
2783                        "keyPackage": key_package,
2784                    }),
2785                )
2786                .await?;
2787            }
2788            ManagedGroupAction::RequestPreparePage {
2789                preparation_id,
2790                page_index,
2791            } => {
2792                execute_operation(
2793                    adapter,
2794                    socket,
2795                    serde_json::json!({
2796                        "v": WIRE_PROTOCOL_VERSION,
2797                        "type": "managed.prepare.page.request",
2798                        "idempotencyKey": random_id("managed-prepare-page"),
2799                        "preparationId": preparation_id,
2800                        "pageIndex": page_index,
2801                    }),
2802                )
2803                .await?;
2804            }
2805            ManagedGroupAction::SendArtifactChunk { chunk } => {
2806                let mut frame = serde_json::to_value(chunk)?;
2807                frame["v"] = serde_json::json!(WIRE_PROTOCOL_VERSION);
2808                frame["type"] = serde_json::json!("managed.commit.chunk");
2809                frame["idempotencyKey"] = serde_json::json!(random_id("managed-artifact"));
2810                execute_operation(adapter, socket, frame).await?;
2811            }
2812            ManagedGroupAction::AcknowledgeCommit {
2813                preparation_id,
2814                encryption_epoch,
2815                group_id_hash,
2816            } => {
2817                execute_operation(
2818                    adapter,
2819                    socket,
2820                    serde_json::json!({
2821                        "v": WIRE_PROTOCOL_VERSION,
2822                        "type": "managed.commit.ack",
2823                        "idempotencyKey": random_id("managed-commit-ack"),
2824                        "preparationId": preparation_id,
2825                        "encryptionEpoch": encryption_epoch,
2826                        "groupIdHash": group_id_hash,
2827                    }),
2828                )
2829                .await?;
2830            }
2831        }
2832    }
2833    Ok(())
2834}
2835
2836#[cfg(feature = "managed-group-encryption")]
2837async fn ensure_managed_group_published(
2838    adapter: &NativeCoordinationGatewaySignaling,
2839    socket: &mut GatewaySocket,
2840) -> Result<()> {
2841    let Some(action) = ensure_managed_group(adapter).await? else {
2842        return Ok(());
2843    };
2844    if let Err(error) = execute_managed_actions(adapter, socket, vec![action]).await {
2845        *adapter.shared.managed_group.lock().await = None;
2846        return Err(error).context("publish native managed group KeyPackage");
2847    }
2848    Ok(())
2849}
2850
2851async fn handle_frame_with_socket(
2852    adapter: &NativeCoordinationGatewaySignaling,
2853    socket: &mut GatewaySocket,
2854    frame: ServerFrame,
2855) -> Result<()> {
2856    #[cfg(feature = "managed-group-encryption")]
2857    match frame {
2858        ServerFrame::TopologyLease { lease } => {
2859            accept_topology_lease(adapter, lease).await?;
2860            if adapter
2861                .room_architecture()
2862                .is_some_and(|snapshot| snapshot.effective == EffectiveRoomArchitecture::Managed)
2863            {
2864                ensure_managed_group_published(adapter, socket).await?;
2865            }
2866            return Ok(());
2867        }
2868        ServerFrame::ManagedPreparePage { page } => {
2869            ensure_managed_group_published(adapter, socket).await?;
2870            let actions = adapter
2871                .shared
2872                .managed_group
2873                .lock()
2874                .await
2875                .as_mut()
2876                .ok_or_else(|| anyhow!("native managed group is unavailable"))?
2877                .handle_prepare_page(page)?;
2878            execute_managed_actions(adapter, socket, actions).await?;
2879            return Ok(());
2880        }
2881        ServerFrame::ManagedArtifactChunk { chunk } => {
2882            ensure_managed_group_published(adapter, socket).await?;
2883            let actions = adapter
2884                .shared
2885                .managed_group
2886                .lock()
2887                .await
2888                .as_mut()
2889                .ok_or_else(|| anyhow!("native managed group is unavailable"))?
2890                .handle_artifact_chunk(chunk)?;
2891            execute_managed_actions(adapter, socket, actions).await?;
2892            return Ok(());
2893        }
2894        ServerFrame::ManagedReceived {
2895            sender_device_id,
2896            architecture_epoch,
2897            encryption_epoch,
2898            batch,
2899        } => {
2900            handle_native_managed_delivery(
2901                adapter,
2902                sender_device_id,
2903                architecture_epoch,
2904                encryption_epoch,
2905                batch,
2906            )
2907            .await?;
2908            return Ok(());
2909        }
2910        other => return handle_frame(adapter, other).await,
2911    }
2912    #[cfg(not(feature = "managed-group-encryption"))]
2913    {
2914        let _ = socket;
2915        handle_frame(adapter, frame).await
2916    }
2917}
2918
2919fn valid_reservation_id(value: &str) -> bool {
2920    !value.is_empty()
2921        && value.len() <= 160
2922        && value.bytes().all(|byte| {
2923            byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b':' | b'@' | b'.')
2924        })
2925}
2926
2927async fn accept_topology_lease(
2928    adapter: &NativeCoordinationGatewaySignaling,
2929    lease: GatewayTopologyLease,
2930) -> Result<()> {
2931    let invalid = || gateway_connect_error("native coordination topology lease is invalid", false);
2932    let current_revision = *adapter.shared.topology_revision.read().await;
2933    let active_grant = adapter.shared.active_grant.read().await.clone();
2934    let Some((grant_jti, grant_expires_at_ms)) = active_grant else {
2935        return Err(invalid());
2936    };
2937    if lease.schema_version != 2
2938        || lease.topology_revision == 0
2939        || lease.avenue != adapter.avenue
2940        || lease.grant_jti != grant_jti
2941        || lease.expires_at_ms <= now_ms()
2942        || lease.expires_at_ms > grant_expires_at_ms
2943    {
2944        return Err(invalid());
2945    }
2946    if lease.topology_revision <= current_revision {
2947        return Ok(());
2948    }
2949    let current_architecture = adapter.shared.room_architecture.borrow().clone();
2950    let current_encryption_epoch = *adapter.shared.encryption_epoch.read().await;
2951    let routes = lease
2952        .active
2953        .iter()
2954        .chain(&lease.backups)
2955        .collect::<Vec<_>>();
2956    let route_ids = routes
2957        .iter()
2958        .map(|route| route.device_id.as_str())
2959        .collect::<std::collections::HashSet<_>>();
2960    let members = adapter.shared.members.read().await;
2961    let current_epoch = current_architecture
2962        .as_ref()
2963        .map_or(0, |snapshot| snapshot.epoch);
2964    let architecture = &lease.architecture;
2965    let valid_peer_admission = architecture.effective_mode != EffectiveRoomArchitecture::Authority
2966        || lease.admission_peers.as_ref().is_some_and(|peers| {
2967            peers.len() <= 50
2968                && peers
2969                    .iter()
2970                    .map(|peer| &peer.device_id)
2971                    .collect::<std::collections::HashSet<_>>()
2972                    .len()
2973                    == peers.len()
2974                && peers.iter().all(|peer| {
2975                    peer.device_id != adapter.device_id
2976                        && members
2977                            .get(&peer.device_id)
2978                            .is_some_and(|member| member.online)
2979                        && peer.node_id.parse::<iroh::EndpointId>().is_ok()
2980                })
2981        });
2982    let valid_group_encryption = match lease.group_encryption.as_ref() {
2983        None => {
2984            architecture.effective_mode != EffectiveRoomArchitecture::Managed
2985                || architecture.phase == RoomArchitecturePhase::Preparing
2986        }
2987        Some(encryption) => {
2988            encryption.policy_version == "room-mls-v1"
2989                && encryption.epoch > 0
2990                && encryption.epoch >= current_encryption_epoch
2991                && (encryption.epoch == current_encryption_epoch
2992                    || current_encryption_epoch == 0
2993                    || encryption.previous_encryption_epoch == Some(current_encryption_epoch))
2994                && matches!(encryption.phase.as_str(), "preparing" | "settled")
2995                && (architecture.effective_mode != EffectiveRoomArchitecture::Managed
2996                    || architecture.phase != RoomArchitecturePhase::Settled
2997                    || encryption.phase == "settled")
2998                && valid_reservation_id(&encryption.committer_device_id)
2999                && encryption.group_id_hash.len() == 43
3000                && encryption
3001                    .group_id_hash
3002                    .bytes()
3003                    .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
3004        }
3005    };
3006    let fixed_request_matches = adapter.architecture.is_none_or(|requested| {
3007        requested == RoomArchitectureMode::Auto || requested == architecture.requested_mode
3008    });
3009    if (current_revision == 0 && lease.previous_topology_revision.is_some())
3010        || (current_revision > 0 && lease.previous_topology_revision != Some(current_revision))
3011        || lease.active.len() > 16
3012        || lease.backups.len() > 4
3013        || routes.len() > 20
3014        || route_ids.len() != routes.len()
3015        || routes.iter().any(|route| {
3016            !route.online
3017                || route.device_id == adapter.device_id
3018                || route.node_id.trim().is_empty()
3019                || route.ticket.trim().is_empty()
3020                || !members.contains_key(&route.device_id)
3021        })
3022        || architecture.policy_version != "room-architecture-v1"
3023        || architecture.epoch == 0
3024        || architecture.epoch < current_epoch
3025        || (architecture.epoch > current_epoch
3026            && current_epoch > 0
3027            && architecture.previous_architecture_epoch != Some(current_epoch))
3028        || !fixed_request_matches
3029        || architecture.held_credits_microusd > 9_007_199_254_740_991
3030        || architecture.quote_expires_at_ms == 0
3031        || architecture
3032            .reservation_id
3033            .as_deref()
3034            .is_some_and(|value| !valid_reservation_id(value))
3035        || !valid_group_encryption
3036        || !valid_peer_admission
3037    {
3038        return Err(invalid());
3039    }
3040    drop(members);
3041
3042    if adapter.authority_client.lock().await.is_some() {
3043        let scope = format!("v2:room:{}", adapter.avenue.id);
3044        if routes.iter().any(|route| {
3045            let (ticket, suffix) = crate::session_token::split_ticket(&route.ticket);
3046            !suffix
3047                .and_then(|suffix| crate::session_token::decode_payload(ticket, suffix))
3048                .is_some_and(|payload| payload.scope.as_str() == scope)
3049        }) {
3050            return Err(gateway_connect_error(
3051                "native authority route lacks room-scoped admission",
3052                false,
3053            ));
3054        }
3055    }
3056
3057    if architecture.effective_mode == EffectiveRoomArchitecture::Authority {
3058        let mut binding = adapter.authority_client.lock().await;
3059        if let Some(binding) = binding.as_mut() {
3060            let client = binding
3061                .client
3062                .upgrade()
3063                .filter(|_| !binding.closed)
3064                .ok_or_else(invalid)?;
3065            binding.admission_revision = binding
3066                .admission_revision
3067                .checked_add(1)
3068                .ok_or_else(invalid)?;
3069            client
3070                .session_token_registry
3071                .update_scope_peer_admission(
3072                    &binding.scope,
3073                    binding.admission_revision,
3074                    lease.expires_at_ms,
3075                    lease
3076                        .admission_peers
3077                        .as_ref()
3078                        .ok_or_else(invalid)?
3079                        .iter()
3080                        .map(|peer| peer.node_id.clone())
3081                        .collect(),
3082                )
3083                .map_err(anyhow::Error::msg)?;
3084        }
3085    }
3086
3087    let snapshot = RoomArchitectureSnapshot {
3088        requested: architecture.requested_mode,
3089        effective: architecture.effective_mode,
3090        epoch: architecture.epoch,
3091        phase: architecture.phase,
3092        reason: architecture.reason,
3093        held_credits_usd: architecture.held_credits_microusd as f64 / 1_000_000.0,
3094        quote_expires_at_ms: architecture.quote_expires_at_ms,
3095    };
3096    let retains_authority_assignment = snapshot.effective == EffectiveRoomArchitecture::Authority;
3097    adapter
3098        .shared
3099        .room_architecture
3100        .send_if_modified(|current| {
3101            if current.as_ref() == Some(&snapshot) {
3102                false
3103            } else {
3104                *current = Some(snapshot);
3105                true
3106            }
3107        });
3108    if !retains_authority_assignment {
3109        adapter.shared.authority_assignment.send_replace(None);
3110    }
3111    if let Some(encryption) = lease.group_encryption.as_ref() {
3112        *adapter.shared.encryption_epoch.write().await = encryption.epoch;
3113    }
3114    *adapter.shared.topology_routes.write().await = lease
3115        .active
3116        .into_iter()
3117        .map(|route| (route, true))
3118        .chain(lease.backups.into_iter().map(|route| (route, false)))
3119        .map(|(route, active)| (route.device_id.clone(), (route, active)))
3120        .collect();
3121    *adapter.shared.topology_revision.write().await = lease.topology_revision;
3122    rebuild_sparse_devices(adapter).await;
3123    adapter.forward_authority_peers().await
3124}
3125
3126fn decode_frame(message: Message) -> Result<Option<ServerFrame>> {
3127    match message {
3128        Message::Text(text) if text == "pong" => Ok(None),
3129        Message::Text(text) => Ok(Some(serde_json::from_str(&text)?)),
3130        Message::Binary(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
3131        Message::Ping(_) | Message::Pong(_) => Ok(None),
3132        Message::Close(Some(frame)) if u16::from(frame.code) == 4403 => Err(gateway_server_error(
3133            "credential-revoked".into(),
3134            "gateway authorization withdrawn".into(),
3135            false,
3136        )),
3137        Message::Close(_) => bail!("coordination gateway closed"),
3138        _ => Ok(None),
3139    }
3140}
3141
3142fn gateway_grant_jti(token: &str) -> Result<String> {
3143    let payload = token
3144        .split('.')
3145        .nth(1)
3146        .ok_or_else(|| anyhow!("native gateway grant is malformed"))?;
3147    let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
3148        .decode(payload)
3149        .context("decode native gateway grant")?;
3150    let value: serde_json::Value = serde_json::from_slice(&decoded)?;
3151    let jti = value
3152        .get("jti")
3153        .and_then(serde_json::Value::as_str)
3154        .unwrap_or_default();
3155    if jti.is_empty()
3156        || jti.len() > 200
3157        || !jti.bytes().all(|byte| {
3158            byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-' | b'.' | b':' | b'@')
3159        })
3160    {
3161        bail!("native gateway grant jti is invalid");
3162    }
3163    Ok(jti.to_string())
3164}
3165
3166async fn handle_frame(
3167    adapter: &NativeCoordinationGatewaySignaling,
3168    frame: ServerFrame,
3169) -> Result<()> {
3170    match frame {
3171        ServerFrame::RosterSnapshot { devices } => {
3172            let mut mapped = HashMap::new();
3173            let mut events = Vec::with_capacity(devices.len());
3174            for raw in devices {
3175                let device = adapter.device_from_gateway(raw);
3176                mapped.insert(device.device_id.clone(), device.clone());
3177                events.push(DeviceEvent::Added { device });
3178            }
3179            *adapter.shared.devices.write().await = mapped;
3180            if !events.is_empty() {
3181                let _ = adapter.shared.device_events.send(events);
3182            }
3183        }
3184        ServerFrame::MembershipSnapshot { members } => {
3185            *adapter
3186                .shared
3187                .membership_pages
3188                .lock()
3189                .map_err(|_| anyhow!("native membership page lock poisoned"))? = None;
3190            *adapter.shared.members.write().await = members
3191                .into_iter()
3192                .map(|member| (member.device_id.clone(), member))
3193                .collect();
3194            rebuild_sparse_devices(adapter).await;
3195        }
3196        ServerFrame::MembershipPage {
3197            snapshot_id,
3198            page_index,
3199            page_count,
3200            member_count,
3201            members,
3202        } => {
3203            let valid_snapshot_id = !snapshot_id.is_empty()
3204                && snapshot_id.len() <= 160
3205                && snapshot_id
3206                    .bytes()
3207                    .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b':' | b'_' | b'-'));
3208            let valid_envelope = valid_snapshot_id
3209                && page_count >= 2
3210                && page_count <= 50
3211                && page_index < page_count
3212                && member_count > 100
3213                && member_count <= 5_000
3214                && !members.is_empty()
3215                && members.len() <= 100;
3216            let completed = {
3217                let mut pending = adapter
3218                    .shared
3219                    .membership_pages
3220                    .lock()
3221                    .map_err(|_| anyhow!("native membership page lock poisoned"))?;
3222                if page_index == 0 && valid_envelope {
3223                    *pending = Some(PendingMembershipPages {
3224                        snapshot_id: snapshot_id.clone(),
3225                        next_page_index: 0,
3226                        page_count,
3227                        member_count,
3228                        members: HashMap::new(),
3229                    });
3230                }
3231                let page = pending
3232                    .as_mut()
3233                    .ok_or_else(|| anyhow!("native membership page arrived without a snapshot"))?;
3234                if !valid_envelope
3235                    || page.snapshot_id != snapshot_id
3236                    || page.next_page_index != page_index
3237                    || page.page_count != page_count
3238                    || page.member_count != member_count
3239                    || members.iter().any(|member| {
3240                        member.device_id.is_empty() || page.members.contains_key(&member.device_id)
3241                    })
3242                {
3243                    bail!("native membership page is invalid");
3244                }
3245                for member in members {
3246                    page.members.insert(member.device_id.clone(), member);
3247                }
3248                page.next_page_index += 1;
3249                if page.next_page_index == page.page_count {
3250                    if page.members.len() != page.member_count {
3251                        bail!("native membership pages are incomplete");
3252                    }
3253                    pending.take().map(|complete| complete.members)
3254                } else {
3255                    None
3256                }
3257            };
3258            if let Some(members) = completed {
3259                *adapter.shared.members.write().await = members;
3260                rebuild_sparse_devices(adapter).await;
3261            }
3262        }
3263        ServerFrame::MembershipChanged { operation, member } => {
3264            let mut members = adapter.shared.members.write().await;
3265            if operation == "delete" {
3266                members.remove(&member.device_id);
3267            } else {
3268                members.insert(member.device_id.clone(), member);
3269            }
3270            drop(members);
3271            rebuild_sparse_devices(adapter).await;
3272        }
3273        ServerFrame::TopologyLease { lease } => {
3274            accept_topology_lease(adapter, lease).await?;
3275        }
3276        ServerFrame::PresenceChanged { operation, device } => {
3277            let device = adapter.device_from_gateway(device);
3278            let event = if operation == "delete" {
3279                adapter
3280                    .shared
3281                    .devices
3282                    .write()
3283                    .await
3284                    .remove(&device.device_id);
3285                DeviceEvent::Removed {
3286                    device_id: device.device_id,
3287                }
3288            } else {
3289                let existed = adapter
3290                    .shared
3291                    .devices
3292                    .write()
3293                    .await
3294                    .insert(device.device_id.clone(), device.clone())
3295                    .is_some();
3296                if existed {
3297                    DeviceEvent::Modified { device }
3298                } else {
3299                    DeviceEvent::Added { device }
3300                }
3301            };
3302            let _ = adapter.shared.device_events.send(vec![event]);
3303        }
3304        ServerFrame::SessionChanged {
3305            operation,
3306            session_id,
3307            session,
3308        } => {
3309            let event = if operation == "delete" {
3310                SessionEvent::Removed { session_id }
3311            } else {
3312                let mut value = session.unwrap_or_else(|| serde_json::json!({}));
3313                value["connectionId"] = serde_json::Value::String(session_id);
3314                SessionEvent::Modified {
3315                    session: serde_json::from_value(value)?,
3316                }
3317            };
3318            let _ = adapter.shared.session_events.send(vec![event]);
3319        }
3320        ServerFrame::SignalReceived {
3321            sender_device_id,
3322            payload,
3323            state,
3324            reply_payload,
3325            ..
3326        } => {
3327            let mut pending = adapter
3328                .shared
3329                .pending_messages
3330                .lock()
3331                .unwrap_or_else(|poisoned| poisoned.into_inner());
3332            if pending.len() >= 1_000 {
3333                pending.remove(0);
3334            }
3335            pending.push(SignalingEnvelope {
3336                app_tag: Some(adapter.app_tag.clone()),
3337                sender_id: sender_device_id,
3338                target_id: adapter.device_id.clone(),
3339                payload,
3340                state,
3341                reply_payload,
3342                timestamp: now_ms() as i64,
3343                sender_user_id: None,
3344                target_user_id: None,
3345                expires_at: None,
3346            });
3347        }
3348        ServerFrame::Error {
3349            code,
3350            message,
3351            retryable,
3352            metadata,
3353            ..
3354        } => return Err(gateway_server_error_metadata(code, message, retryable, metadata)),
3355        ServerFrame::AuthorityAssignment {
3356            assignment,
3357            signature,
3358        } => {
3359            validate_authority_assignment(&assignment, &signature)?;
3360            if assignment.subject_device_id != adapter.device_id
3361                || !matches!(
3362                    adapter.shared.room_architecture.borrow().as_ref(),
3363                    Some(snapshot) if snapshot.effective == EffectiveRoomArchitecture::Authority
3364                )
3365            {
3366                bail!("native authority assignment does not match the active room lease");
3367            }
3368            let next = NativeAuthorityAssignment {
3369                assignment,
3370                signature,
3371            };
3372            if let Some(current) = adapter.shared.authority_assignment.borrow().clone() {
3373                let current_key = (current.assignment.generation, current.assignment.revision);
3374                let next_key = (next.assignment.generation, next.assignment.revision);
3375                if next_key < current_key || (next_key == current_key && next != current) {
3376                    bail!("native authority assignment generation is stale");
3377                }
3378                if next == current {
3379                    return Ok(());
3380                }
3381            }
3382            adapter.shared.authority_assignment.send_replace(Some(next));
3383        }
3384        #[cfg(feature = "managed-group-encryption")]
3385        ServerFrame::ManagedPreparePage { .. }
3386        | ServerFrame::ManagedArtifactChunk { .. }
3387        | ServerFrame::ManagedReceived { .. } => {
3388            bail!("native managed room frame requires the active gateway owner")
3389        }
3390        ServerFrame::Ready { .. }
3391        | ServerFrame::AuthRefreshed { .. }
3392        | ServerFrame::LeaseRefreshed { .. }
3393        | ServerFrame::Ack { .. }
3394        | ServerFrame::Pong => {}
3395    }
3396    adapter.forward_authority_peers().await
3397}
3398
3399async fn presence_frame(
3400    adapter: &NativeCoordinationGatewaySignaling,
3401    desired: &DesiredPresence,
3402) -> serde_json::Value {
3403    serde_json::json!({
3404        "v": WIRE_PROTOCOL_VERSION,
3405        "type": "presence.upsert",
3406        "idempotencyKey": random_id("presence"),
3407        "ttlMs": desired.ttl_ms,
3408        "device": gateway_device_value(adapter, desired).await,
3409    })
3410}
3411
3412async fn gateway_device_value(
3413    adapter: &NativeCoordinationGatewaySignaling,
3414    desired: &DesiredPresence,
3415) -> serde_json::Value {
3416    let patch = adapter.shared.staged_patch.read().await.clone();
3417    serde_json::to_value(OutboundGatewayDevice {
3418        device_id: adapter.device_id.clone(),
3419        runtime_instance_id: adapter.runtime_instance_id.clone(),
3420        node_id: desired.local_node_id.clone(),
3421        device_name: patch
3422            .device_name
3423            .unwrap_or_else(|| desired.device_name.clone()),
3424        platform_type: adapter.platform_type.clone(),
3425        ticket: desired.ticket.clone(),
3426        metadata: patch.metadata.or_else(|| desired.metadata.clone()),
3427        capabilities: patch.capabilities,
3428        excluded_peers: patch.excluded_peers.unwrap_or_default(),
3429        online: desired.online,
3430    })
3431    .expect("gateway device projection is serializable")
3432}
3433
3434fn merge_patch(target: &mut DevicePatch, patch: DevicePatch) {
3435    if patch.device_name.is_some() {
3436        target.device_name = patch.device_name;
3437    }
3438    if patch.capabilities.is_some() {
3439        target.capabilities = patch.capabilities;
3440    }
3441    if patch.metadata.is_some() {
3442        target.metadata = patch.metadata;
3443    }
3444    if patch.excluded_peers.is_some() {
3445        target.excluded_peers = patch.excluded_peers;
3446    }
3447}
3448
3449fn reject_command(command: Command, message: &str) {
3450    match command {
3451        Command::Publish { reply, .. }
3452        | Command::Patch { reply, .. }
3453        | Command::Offline { reply }
3454        | Command::Delete { reply, .. }
3455        | Command::PutSession { reply, .. } => {
3456            let _ = reply.send(Err(anyhow!(message.to_string())));
3457        }
3458        Command::PublishAuthorityAssignment { reply, .. } => {
3459            let _ = reply.send(Err(anyhow!(message.to_string())));
3460        }
3461        #[cfg(feature = "managed-group-encryption")]
3462        Command::PublishManaged { reply, .. } => {
3463            let _ = reply.send(Err(anyhow!(message.to_string())));
3464        }
3465        Command::SendSignal { reply, .. } => {
3466            let _ = reply.send(Err(anyhow!(message.to_string())));
3467        }
3468        Command::Stop => {}
3469    }
3470}
3471
3472#[cfg(feature = "managed-group-encryption")]
3473async fn publish_native_managed_batch(
3474    adapter: &NativeCoordinationGatewaySignaling,
3475    socket: &mut GatewaySocket,
3476    batch: Vec<NativeManagedRoomPublish>,
3477) -> Result<()> {
3478    if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
3479        bail!("managed room publish batch exceeds its protected bound");
3480    }
3481    let architecture = adapter
3482        .room_architecture()
3483        .ok_or_else(|| anyhow!("native managed room architecture is unavailable"))?;
3484    let encryption_epoch = *adapter.shared.encryption_epoch.read().await;
3485    if architecture.effective != EffectiveRoomArchitecture::Managed
3486        || architecture.phase != RoomArchitecturePhase::Settled
3487        || encryption_epoch == 0
3488    {
3489        bail!("native managed room encryption is not settled");
3490    }
3491    let (envelopes, sealed_state) = {
3492        let mut managed = adapter.shared.managed_group.lock().await;
3493        let controller = managed
3494            .as_mut()
3495            .ok_or_else(|| anyhow!("native managed room group is unavailable"))?;
3496        let mut envelopes = Vec::with_capacity(batch.len());
3497        let mut sealed_state = None;
3498        for message in batch {
3499            let protected = controller.seal_payload(
3500                architecture.epoch,
3501                encryption_epoch,
3502                &message.message_id,
3503                &message.channel,
3504                message.priority,
3505                message.zone_id.as_deref(),
3506                &message.payload,
3507            )?;
3508            sealed_state = Some(protected.sealed_state);
3509            envelopes.push(NativeManagedEncryptedEnvelope {
3510                message_id: message.message_id,
3511                channel: message.channel,
3512                priority: message.priority,
3513                scope: message
3514                    .zone_id
3515                    .map_or(NativeManagedScope::Global, |zone_id| {
3516                        NativeManagedScope::Zone { zone_id }
3517                    }),
3518                ciphertext: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(protected.data),
3519            });
3520        }
3521        (
3522            envelopes,
3523            sealed_state.expect("non-empty managed batch has sealed state"),
3524        )
3525    };
3526    adapter
3527        .managed_group_store
3528        .as_ref()
3529        .ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
3530        .save(
3531            &adapter.app_tag,
3532            &adapter.device_id,
3533            &adapter.avenue_key(),
3534            &sealed_state,
3535        )
3536        .await
3537        .context("persist native managed group state before publish")?;
3538    execute_operation(
3539        adapter,
3540        socket,
3541        serde_json::json!({
3542            "v": WIRE_PROTOCOL_VERSION,
3543            "type": "managed.publish",
3544            "idempotencyKey": random_id("managed-publish"),
3545            "architectureEpoch": architecture.epoch,
3546            "encryptionEpoch": encryption_epoch,
3547            "batch": envelopes,
3548        }),
3549    )
3550    .await
3551}
3552
3553#[cfg(feature = "managed-group-encryption")]
3554async fn handle_native_managed_delivery(
3555    adapter: &NativeCoordinationGatewaySignaling,
3556    sender_device_id: String,
3557    architecture_epoch: u64,
3558    encryption_epoch: u64,
3559    batch: Vec<NativeManagedEncryptedEnvelope>,
3560) -> Result<()> {
3561    if batch.is_empty() || batch.len() > MAX_MANAGED_BATCH_MESSAGES {
3562        bail!("native managed delivery batch exceeds its protected bound");
3563    }
3564    let architecture = adapter
3565        .room_architecture()
3566        .ok_or_else(|| anyhow!("native managed room architecture is unavailable"))?;
3567    if architecture.effective != EffectiveRoomArchitecture::Managed
3568        || architecture.phase != RoomArchitecturePhase::Settled
3569        || architecture.epoch != architecture_epoch
3570        || *adapter.shared.encryption_epoch.read().await != encryption_epoch
3571    {
3572        bail!("native managed room delivery uses a stale lease");
3573    }
3574    let (messages, sealed_state) = {
3575        let mut managed = adapter.shared.managed_group.lock().await;
3576        let controller = managed
3577            .as_mut()
3578            .ok_or_else(|| anyhow!("native managed room group is unavailable"))?;
3579        let mut messages = Vec::with_capacity(batch.len());
3580        let mut sealed_state = None;
3581        for envelope in batch {
3582            let zone_id = match envelope.scope {
3583                NativeManagedScope::Global => None,
3584                NativeManagedScope::Zone { zone_id } => Some(zone_id),
3585            };
3586            let ciphertext = base64::engine::general_purpose::URL_SAFE_NO_PAD
3587                .decode(envelope.ciphertext)
3588                .context("decode native managed ciphertext")?;
3589            let protected = controller.open_payload(
3590                architecture_epoch,
3591                encryption_epoch,
3592                &envelope.message_id,
3593                &envelope.channel,
3594                envelope.priority,
3595                zone_id.as_deref(),
3596                &ciphertext,
3597            )?;
3598            sealed_state = Some(protected.sealed_state);
3599            messages.push(NativeManagedRoomMessage {
3600                sender_device_id: sender_device_id.clone(),
3601                message_id: envelope.message_id,
3602                channel: envelope.channel,
3603                priority: envelope.priority,
3604                zone_id,
3605                payload: protected.data,
3606            });
3607        }
3608        (
3609            messages,
3610            sealed_state.expect("non-empty managed delivery has sealed state"),
3611        )
3612    };
3613    adapter
3614        .managed_group_store
3615        .as_ref()
3616        .ok_or_else(|| anyhow!("native managed group state store is unavailable"))?
3617        .save(
3618            &adapter.app_tag,
3619            &adapter.device_id,
3620            &adapter.avenue_key(),
3621            &sealed_state,
3622        )
3623        .await
3624        .context("persist native managed group state before delivery")?;
3625    for message in messages {
3626        let _ = adapter.shared.managed_messages.send(message);
3627    }
3628    Ok(())
3629}
3630
3631fn reconnect_delay(attempt: u8) -> Duration {
3632    let exponent = u32::from(attempt.saturating_sub(1).min(7));
3633    Duration::from_millis(
3634        (250_u64.saturating_mul(1_u64 << exponent)).min(MAX_RECONNECT_DELAY.as_millis() as u64),
3635    )
3636}
3637
3638fn gateway_reconnect_delay(
3639    attempt: u8,
3640    error: &anyhow::Error,
3641    authority_return_retry_available: &mut bool,
3642) -> Option<Duration> {
3643    if !is_retryable_gateway_error(error) || attempt >= MAX_RECONNECT_ATTEMPTS {
3644        return None;
3645    }
3646    // Reviewed 2026-09-05; native gateway actor / RoomAuthority consumers.
3647    // After an established socket's outage, an authenticated gateway can return
3648    // just before its host. Spend ONE existing attempt promptly on that typed
3649    // response, without resetting the outage counter or granting admission.
3650    // Only successful authentication rearms this; repeated absence keeps backoff.
3651    if *authority_return_retry_available
3652        && error
3653            .downcast_ref::<GatewayConnectError>()
3654            .is_some_and(|error| error.code.as_deref() == Some("room-authority-unavailable"))
3655    {
3656        *authority_return_retry_available = false;
3657        return Some(reconnect_delay(1));
3658    }
3659    // Financial cooldowns are minimum admission times, not network backoff.
3660    // Never cap them to the 30-second transport backoff ceiling.
3661    Some(reconnect_delay(attempt).max(gateway_service_retry_delay(error)))
3662}
3663
3664fn auth_refresh_delay(expires_at_ms: u64) -> Duration {
3665    Duration::from_millis(
3666        expires_at_ms
3667            .saturating_sub(now_ms())
3668            .saturating_sub(AUTH_REFRESH_SKEW_MS)
3669            .max(1_000),
3670    )
3671}
3672
3673async fn send_socket_keepalive<S, E>(socket: &mut S) -> Result<()>
3674where
3675    S: Sink<Message, Error = E> + Unpin,
3676    E: StdError + Send + Sync + 'static,
3677{
3678    socket
3679        .send(Message::Text("ping".into()))
3680        .await
3681        .context("send native coordination WebSocket keepalive")
3682}
3683
3684fn required(name: &str, value: String) -> Result<String> {
3685    let value = value.trim().to_string();
3686    if value.is_empty() {
3687        bail!("native coordination {name} is required");
3688    }
3689    Ok(value)
3690}
3691
3692fn validate_endpoint(name: &str, value: String, allow_websocket: bool) -> Result<String> {
3693    let value = required(name, value)?;
3694    let url = reqwest::Url::parse(&value)?;
3695    let allowed = matches!(url.scheme(), "https" | "http")
3696        || (allow_websocket && matches!(url.scheme(), "wss" | "ws"));
3697    let secure = matches!(url.scheme(), "https" | "wss");
3698    let loopback = matches!(
3699        url.host_str(),
3700        Some("127.0.0.1") | Some("localhost") | Some("::1")
3701    );
3702    if !allowed
3703        || (!secure && !loopback)
3704        || !url.username().is_empty()
3705        || url.password().is_some()
3706        || url.query().is_some()
3707        || url.fragment().is_some()
3708    {
3709        bail!("native coordination {name} is invalid");
3710    }
3711    Ok(value.trim_end_matches('/').to_string())
3712}
3713
3714fn normalized_origin_scheme(scheme: &str) -> &str {
3715    match scheme {
3716        "ws" => "http",
3717        "wss" => "https",
3718        other => other,
3719    }
3720}
3721
3722fn random_id(prefix: &str) -> String {
3723    format!("{prefix}:{}", Uuid::new_v4())
3724}
3725
3726fn now_ms() -> u64 {
3727    SystemTime::now()
3728        .duration_since(UNIX_EPOCH)
3729        .unwrap_or_default()
3730        .as_millis() as u64
3731}
3732
3733#[async_trait]
3734impl SignalingBackend for NativeCoordinationGatewaySignaling {
3735    async fn update_presence(
3736        &self,
3737        user_id: &str,
3738        local_node_id: &str,
3739        ticket_str: &str,
3740        is_online: bool,
3741        name: &str,
3742        ttl_ms: u64,
3743        metadata: Option<&str>,
3744    ) -> Result<()> {
3745        // The native lifecycle exposes its historical durable-device and
3746        // live-lease projections as two concurrent backend calls. The gateway
3747        // owns both projections in one presence upsert, so retain the durable
3748        // fields locally and let the live call perform the single billed
3749        // operation. Explicit offline transitions use `set_offline`.
3750        if !is_online {
3751            merge_patch(
3752                &mut *self.shared.staged_patch.write().await,
3753                DevicePatch {
3754                    device_name: Some(name.to_string()),
3755                    metadata: metadata.map(str::to_string),
3756                    ..DevicePatch::default()
3757                },
3758            );
3759            return Ok(());
3760        }
3761        let local_node_id = iroh::EndpointId::from_str(local_node_id)
3762            .context("native coordination node ID is invalid")?
3763            .to_string();
3764        let desired = DesiredPresence {
3765            user_id: user_id.to_string(),
3766            local_node_id,
3767            ticket: ticket_str.to_string(),
3768            device_name: name.to_string(),
3769            metadata: metadata.map(str::to_string),
3770            ttl_ms,
3771            online: is_online,
3772        };
3773        *self.shared.desired.write().await = Some(desired.clone());
3774        self.request_unit(|reply| Command::Publish { desired, reply })
3775            .await
3776    }
3777
3778    async fn set_offline(&self, _user_id: &str, _local_node_id: &str) -> Result<()> {
3779        self.request_unit(|reply| Command::Offline { reply }).await
3780    }
3781
3782    async fn update_live_presence(
3783        &self,
3784        user_id: &str,
3785        local_node_id: &str,
3786        ticket_str: &str,
3787        name: &str,
3788        metadata: Option<&str>,
3789    ) -> Result<()> {
3790        self.update_presence(
3791            user_id,
3792            local_node_id,
3793            ticket_str,
3794            true,
3795            name,
3796            15 * 60_000,
3797            metadata,
3798        )
3799        .await
3800    }
3801
3802    async fn set_live_presence_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
3803        self.set_offline(user_id, local_node_id).await
3804    }
3805
3806    async fn update_device(
3807        &self,
3808        _user_id: &str,
3809        device_id: &str,
3810        device_name: Option<&str>,
3811        capabilities: Option<DeviceCapabilities>,
3812        metadata: Option<&str>,
3813    ) -> Result<()> {
3814        if device_id != self.device_id {
3815            bail!("native coordination socket can update only its local device");
3816        }
3817        let patch = DevicePatch {
3818            device_name: device_name.map(str::to_string),
3819            capabilities,
3820            metadata: metadata.map(str::to_string),
3821            excluded_peers: None,
3822        };
3823        merge_patch(&mut *self.shared.staged_patch.write().await, patch.clone());
3824        // Initial durable projection and live lease are folded into the first
3825        // presence upsert. Avoid a second billed device.patch while the
3826        // concurrent native presence actor is still establishing the socket.
3827        if self.shared.applied_presence.read().await.is_none() {
3828            return Ok(());
3829        }
3830        self.request_unit(|reply| Command::Patch { patch, reply })
3831            .await
3832    }
3833
3834    async fn delete_device(&self, _user_id: &str, device_id: &str) -> Result<()> {
3835        self.request_unit(|reply| Command::Delete {
3836            user_id: _user_id.to_string(),
3837            device_id: device_id.to_string(),
3838            reply,
3839        })
3840        .await
3841    }
3842
3843    async fn set_excluded_peers(
3844        &self,
3845        _user_id: &str,
3846        _local_node_id: &str,
3847        excluded_peers: &[String],
3848    ) -> Result<()> {
3849        let mut normalized = excluded_peers
3850            .iter()
3851            .map(|value| value.trim().to_ascii_lowercase())
3852            .filter(|value| !value.is_empty())
3853            .collect::<Vec<_>>();
3854        normalized.sort();
3855        normalized.dedup();
3856        if normalized.len() > MAX_EXCLUDED_PEERS {
3857            bail!(
3858                "native coordination excluded-peer set exceeds {} devices",
3859                MAX_EXCLUDED_PEERS
3860            );
3861        }
3862        let patch = DevicePatch {
3863            excluded_peers: Some(normalized),
3864            ..DevicePatch::default()
3865        };
3866        merge_patch(&mut *self.shared.staged_patch.write().await, patch.clone());
3867        if self.shared.applied_presence.read().await.is_none() {
3868            return Ok(());
3869        }
3870        self.request_unit(|reply| Command::Patch { patch, reply })
3871            .await
3872    }
3873
3874    async fn search_devices(
3875        &self,
3876        _user_id: &str,
3877        exclude_node_id: Option<&str>,
3878    ) -> Result<Vec<Device>> {
3879        self.list_devices("", exclude_node_id).await
3880    }
3881
3882    async fn list_devices(
3883        &self,
3884        _user_id: &str,
3885        exclude_node_id: Option<&str>,
3886    ) -> Result<Vec<Device>> {
3887        let mut devices = self
3888            .shared
3889            .devices
3890            .read()
3891            .await
3892            .values()
3893            .filter(|device| exclude_node_id != device.node_id.as_deref())
3894            .cloned()
3895            .collect::<Vec<_>>();
3896        devices.sort_by(|left, right| left.device_id.cmp(&right.device_id));
3897        Ok(devices)
3898    }
3899
3900    async fn send_message(
3901        &self,
3902        _sender_id: &str,
3903        target_id: &str,
3904        payload: &str,
3905        state: Option<&str>,
3906        reply_payload: Option<&str>,
3907    ) -> Result<String> {
3908        let (reply, response) = oneshot::channel();
3909        self.commands
3910            .send(Command::SendSignal {
3911                target_device_id: target_id.to_string(),
3912                payload: payload.to_string(),
3913                state: state.map(str::to_string),
3914                reply_payload: reply_payload.map(str::to_string),
3915                reply,
3916            })
3917            .await
3918            .map_err(|_| anyhow!("coordination gateway actor stopped"))?;
3919        response
3920            .await
3921            .map_err(|_| anyhow!("coordination gateway actor dropped its reply"))?
3922    }
3923
3924    async fn subscribe_devices(
3925        &self,
3926        _user_id: &str,
3927    ) -> Result<BoxStream<'static, Result<Vec<DeviceEvent>>>> {
3928        let receiver = self.shared.device_events.subscribe();
3929        Ok(Box::pin(
3930            tokio_stream::wrappers::BroadcastStream::new(receiver)
3931                .filter_map(|event| async move { event.ok().map(Ok) }),
3932        ))
3933    }
3934
3935    async fn create_session(&self, session: SignalingSession) -> Result<()> {
3936        let expires_at_ms = session
3937            .expires_at
3938            .unwrap_or_else(|| now_ms() as i64 + 15 * 60_000);
3939        let session_id = session.connection_id.clone();
3940        let value = serde_json::to_value(session)?;
3941        self.request_unit(|reply| Command::PutSession {
3942            session_id,
3943            session: value,
3944            expires_at_ms,
3945            reply,
3946        })
3947        .await
3948    }
3949
3950    async fn update_session(&self, session_id: &str, update_data: serde_json::Value) -> Result<()> {
3951        let expires_at_ms = update_data
3952            .get("expiresAt")
3953            .and_then(serde_json::Value::as_i64)
3954            .unwrap_or_else(|| now_ms() as i64 + 15 * 60_000);
3955        self.request_unit(|reply| Command::PutSession {
3956            session_id: session_id.to_string(),
3957            session: update_data,
3958            expires_at_ms,
3959            reply,
3960        })
3961        .await
3962    }
3963
3964    async fn subscribe_sessions(
3965        &self,
3966        _local_device_id: &str,
3967    ) -> Result<BoxStream<'static, Result<Vec<SessionEvent>>>> {
3968        let receiver = self.shared.session_events.subscribe();
3969        Ok(Box::pin(
3970            tokio_stream::wrappers::BroadcastStream::new(receiver)
3971                .filter_map(|event| async move { event.ok().map(Ok) }),
3972        ))
3973    }
3974}
3975
3976impl Drop for NativeCoordinationGatewaySignaling {
3977    fn drop(&mut self) {
3978        let _ = self.commands.try_send(Command::Stop);
3979    }
3980}
3981
3982#[cfg(test)]
3983mod tests {
3984    use super::*;
3985
3986    #[derive(Default)]
3987    struct RecordingGrantProvider {
3988        requests: Mutex<Vec<NativeGatewayGrantRequest>>,
3989        gateway_url: Option<String>,
3990        token: Option<String>,
3991        blocked: Option<Arc<tokio::sync::Notify>>,
3992    }
3993
3994    #[async_trait]
3995    impl NativeGatewayGrantProvider for RecordingGrantProvider {
3996        async fn grant(&self, request: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
3997            self.requests.lock().expect("requests").push(request);
3998            if let Some(blocked) = &self.blocked { blocked.notified().await; }
3999            Ok(NativeGatewayGrant {
4000                protocol_version: 2,
4001                gateway_url: self
4002                    .gateway_url
4003                    .clone()
4004                    .unwrap_or_else(|| "https://gateway.example.test".into()),
4005                route_key: "route-1".to_string(),
4006                token: self.token.clone().unwrap_or_else(|| "grant-2".into()),
4007                expires_at_ms: now_ms() + 3_600_000,
4008            })
4009        }
4010    }
4011
4012    #[tokio::test]
4013    async fn socket_lease_refresh_updates_only_acknowledged_grant_expiry() {
4014        // Actual loopback socket, fake budget authority. No hosted services.
4015        for outcome in ["accepted", "wrong-request", "denied", "expired"] {
4016            let provider = Arc::new(RecordingGrantProvider::default());
4017            let handle = NativeCapabilities::new(
4018                crate::test_constants::TEST_API_KEY,
4019                "lease-client",
4020                "native",
4021                provider.clone(),
4022            )
4023            .unwrap()
4024            .join_room_with_architecture("lease-room", RoomArchitectureMode::Authority)
4025            .unwrap();
4026            let adapter = &handle.signaling;
4027            let old_expiry = now_ms() + 30_000;
4028            let renewed_expiry = if outcome == "expired" {
4029                1
4030            } else {
4031                old_expiry + 60_000
4032            };
4033            *adapter.shared.active_grant.write().await =
4034                Some(("same-socket-grant".into(), old_expiry));
4035            let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4036            let address = listener.local_addr().unwrap();
4037            let server = tokio::spawn(async move {
4038                let (tcp, _) = listener.accept().await.unwrap();
4039                let mut socket = tokio_tungstenite::accept_async(tcp).await.unwrap();
4040                let request: serde_json::Value =
4041                    serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
4042                        .unwrap();
4043                assert_eq!(request["type"], "lease.refresh");
4044                let response = if outcome == "denied" {
4045                    serde_json::json!({"type":"error", "code":"credential-revoked", "message":"local denial", "retryable":false})
4046                } else {
4047                    serde_json::json!({"type":"lease.refreshed", "expiresAtMs":renewed_expiry,
4048                        "idempotencyKey":if outcome == "wrong-request" { serde_json::json!("other-request") } else { request["idempotencyKey"].clone() }})
4049                };
4050                socket
4051                    .send(Message::Text(response.to_string().into()))
4052                    .await
4053                    .unwrap();
4054                socket.close(None).await.unwrap();
4055            });
4056            let (mut socket, _) = tokio_tungstenite::connect_async(format!("ws://{address}"))
4057                .await
4058                .unwrap();
4059            let result =
4060                tokio::time::timeout(Duration::from_secs(3), refresh_lease(adapter, &mut socket))
4061                    .await
4062                    .unwrap();
4063            let accepted = outcome == "accepted";
4064            assert_eq!(result.is_ok(), accepted, "{outcome}: {result:?}");
4065            assert_eq!(*adapter.shared.active_grant.read().await,
4066                Some(("same-socket-grant".into(), if accepted { renewed_expiry } else { old_expiry })),
4067                "{outcome}: topology validation must use only the acknowledged socket authorization");
4068            assert!(
4069                provider.requests.lock().unwrap().is_empty(),
4070                "lease renewal must not mint a grant"
4071            );
4072            tokio::time::timeout(Duration::from_secs(3), server)
4073                .await
4074                .unwrap()
4075                .unwrap();
4076            handle.close().await;
4077        }
4078    }
4079
4080    #[tokio::test]
4081    async fn authority_ticket_refresh_keeps_socket_and_orders_authorization_before_presence() {
4082        // Real loopback WebSockets, fake grant authority. This proves protocol
4083        // ordering and failure handling, not hosted grant/budget validation.
4084        for outcome in [
4085            "accepted",
4086            "budget-replay",
4087            "auth-denied",
4088            "presence-denied",
4089            "auth-eof",
4090        ] {
4091            let accepted = matches!(outcome, "accepted" | "budget-replay");
4092            let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4093            let endpoint = format!("http://{}", listener.local_addr().unwrap());
4094            let token = format!(
4095                "local.{}.test",
4096                base64::engine::general_purpose::URL_SAFE_NO_PAD
4097                    .encode(br#"{"jti":"renewed-local-grant"}"#)
4098            );
4099            let provider = Arc::new(RecordingGrantProvider {
4100                gateway_url: Some(endpoint.clone()),
4101                token: Some(token.clone()),
4102                ..Default::default()
4103            });
4104            let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
4105                endpoint: endpoint.clone(),
4106                app_tag: "app_native_test".into(),
4107                device_id: "service-device".into(),
4108                platform_type: "native".into(),
4109                avenue: NativeCoordinationAvenue {
4110                    kind: "room".into(),
4111                    id: "local-room".into(),
4112                },
4113                architecture: Some(RoomArchitectureMode::Authority),
4114                room_delivery: None,
4115                grant_provider: provider.clone(),
4116                #[cfg(feature = "managed-group-encryption")]
4117                managed_group_signer: None,
4118                #[cfg(feature = "managed-group-encryption")]
4119                managed_group_store: None,
4120            })
4121            .unwrap();
4122            let old = DesiredPresence {
4123                user_id: "local-room".into(),
4124                local_node_id: "local-node".into(),
4125                ticket: "old-local-ticket".into(),
4126                device_name: "Service".into(),
4127                metadata: None,
4128                ttl_ms: 900_000,
4129                online: true,
4130            };
4131            let desired = DesiredPresence {
4132                ticket: "renewed-local-ticket".into(),
4133                ..old.clone()
4134            };
4135            *adapter.shared.desired.write().await = Some(desired.clone());
4136            *adapter.shared.applied_presence.write().await = Some(old.clone());
4137            let server_adapter = adapter.clone();
4138            let server = tokio::spawn(async move {
4139                let (tcp, _) = listener.accept().await.unwrap();
4140                let mut socket = tokio_tungstenite::accept_async(tcp).await.unwrap();
4141                let frame: serde_json::Value =
4142                    serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
4143                        .unwrap();
4144                assert_eq!(frame["type"], "auth.refresh");
4145                assert_eq!(
4146                    server_adapter
4147                        .shared
4148                        .applied_presence
4149                        .read()
4150                        .await
4151                        .as_ref()
4152                        .unwrap()
4153                        .ticket,
4154                    "old-local-ticket"
4155                );
4156                if outcome == "auth-eof" {
4157                    socket.close(None).await.unwrap();
4158                    return;
4159                }
4160                let response = if outcome == "auth-denied" {
4161                    serde_json::json!({"type":"error", "code":"credential-revoked", "message":"local denial", "retryable":false})
4162                } else {
4163                    // A concurrent caller may stage another desired snapshot.
4164                    server_adapter
4165                        .shared
4166                        .desired
4167                        .write()
4168                        .await
4169                        .as_mut()
4170                        .unwrap()
4171                        .ticket = "later-local-ticket".into();
4172                    serde_json::json!({"type":"auth.refreshed", "expiresAtMs":now_ms()+3_600_000})
4173                };
4174                socket
4175                    .send(Message::Text(response.to_string().into()))
4176                    .await
4177                    .unwrap();
4178                let next = socket.next().await.unwrap().unwrap();
4179                if outcome == "auth-denied" {
4180                    assert!(
4181                        matches!(next, Message::Close(_)),
4182                        "denied auth must not publish"
4183                    );
4184                    return;
4185                }
4186                let frame: serde_json::Value =
4187                    serde_json::from_str(next.to_text().unwrap()).unwrap();
4188                assert_eq!(frame["type"], "presence.upsert");
4189                assert_eq!(frame["device"]["ticket"], "renewed-local-ticket");
4190                if outcome == "budget-replay" {
4191                    socket
4192                        .send(Message::Text(
4193                            serde_json::json!({
4194                                "type":"error", "code":"budget-renewal-required",
4195                                "message":"local renewal", "retryable":true,
4196                                "idempotencyKey":frame["idempotencyKey"],
4197                            })
4198                            .to_string()
4199                            .into(),
4200                        ))
4201                        .await
4202                        .unwrap();
4203                    let refresh: serde_json::Value = serde_json::from_str(
4204                        socket.next().await.unwrap().unwrap().to_text().unwrap(),
4205                    )
4206                    .unwrap();
4207                    assert_eq!(refresh["type"], "auth.refresh");
4208                    socket.send(Message::Text(serde_json::json!({"type":"auth.refreshed", "expiresAtMs":now_ms()+3_600_000}).to_string().into())).await.unwrap();
4209                    let replay: serde_json::Value = serde_json::from_str(
4210                        socket.next().await.unwrap().unwrap().to_text().unwrap(),
4211                    )
4212                    .unwrap();
4213                    assert_eq!(replay, frame, "replay keeps the ticket and idempotency key");
4214                }
4215                assert_eq!(
4216                    server_adapter
4217                        .shared
4218                        .applied_presence
4219                        .read()
4220                        .await
4221                        .as_ref()
4222                        .unwrap()
4223                        .ticket,
4224                    "old-local-ticket"
4225                );
4226                let response = if outcome == "presence-denied" {
4227                    serde_json::json!({"type":"error", "code":"permission-denied", "message":"local denial", "retryable":false, "idempotencyKey":frame["idempotencyKey"]})
4228                } else {
4229                    serde_json::json!({"type":"ack", "idempotencyKey":frame["idempotencyKey"]})
4230                };
4231                socket
4232                    .send(Message::Text(response.to_string().into()))
4233                    .await
4234                    .unwrap();
4235                let next = socket.next().await.unwrap().unwrap();
4236                if accepted {
4237                    assert!(
4238                        matches!(next, Message::Text(payload) if payload == "ping"),
4239                        "successful refresh keeps the same socket"
4240                    );
4241                } else {
4242                    assert!(matches!(next, Message::Close(_)));
4243                }
4244            });
4245            let (socket, _) = tokio_tungstenite::connect_async(endpoint.replace("http:", "ws:"))
4246                .await
4247                .unwrap();
4248            let mut active = ConnectedGateway {
4249                socket,
4250                credential_token: "previous-local-grant".into(),
4251                credential_expires_at_ms: now_ms() + 30_000,
4252                lease_refresh_mode: LeaseRefreshMode::Active,
4253            };
4254            let result = tokio::time::timeout(
4255                Duration::from_secs(3),
4256                refresh_presence_authentication(&adapter, &mut active, desired.clone()),
4257            )
4258            .await
4259            .unwrap();
4260            assert_eq!(result.is_ok(), accepted, "{outcome}: {result:?}");
4261            if accepted {
4262                assert_eq!(*adapter.shared.applied_presence.read().await, Some(desired));
4263                assert_eq!(active.credential_token, token);
4264                assert!(active.credential_expires_at_ms > now_ms() + 30_000);
4265                send_socket_keepalive(&mut active.socket).await.unwrap();
4266            } else {
4267                assert_eq!(*adapter.shared.applied_presence.read().await, Some(old));
4268                if outcome != "auth-eof" {
4269                    assert!(!is_retryable_gateway_error(result.as_ref().unwrap_err()));
4270                    active.socket.close(None).await.unwrap();
4271                }
4272            }
4273            tokio::time::timeout(Duration::from_secs(3), server)
4274                .await
4275                .unwrap()
4276                .unwrap();
4277            let requests = provider.requests.lock().unwrap();
4278            assert_eq!(
4279                requests.len(),
4280                if outcome == "budget-replay" { 2 } else { 1 },
4281                "only the existing bounded budget replay may request another grant"
4282            );
4283            if outcome == "budget-replay" {
4284                assert_eq!(requests[1].refresh_grant, None);
4285                assert_eq!(
4286                    requests[1].ticket_fingerprint,
4287                    requests[0].ticket_fingerprint
4288                );
4289            }
4290            assert_eq!(
4291                requests[0].refresh_grant.as_deref(),
4292                Some("previous-local-grant")
4293            );
4294            assert_eq!(
4295                requests[0].ticket_fingerprint,
4296                base64::engine::general_purpose::URL_SAFE_NO_PAD
4297                    .encode(Sha256::digest(b"renewed-local-ticket"))
4298            );
4299            drop(requests);
4300            adapter.stop().await;
4301        }
4302    }
4303
4304    #[test]
4305    fn native_authority_targets_match_the_public_assignment_bounds() {
4306        let priority = (0..8)
4307            .map(|index| format!("peer-{index}"))
4308            .collect::<Vec<_>>();
4309        let relevant = (0..32)
4310            .map(|index| format!("entity-{index}"))
4311            .collect::<Vec<_>>();
4312        validate_native_authority_assignment_targets(&priority, &relevant)
4313            .expect("the exact public bounds are accepted");
4314
4315        let mut too_many_priority = priority.clone();
4316        too_many_priority.push("peer-8".to_string());
4317        assert!(
4318            validate_native_authority_assignment_targets(&too_many_priority, &relevant).is_err()
4319        );
4320
4321        let mut too_many_relevant = relevant.clone();
4322        too_many_relevant.push("entity-32".to_string());
4323        assert!(
4324            validate_native_authority_assignment_targets(&priority, &too_many_relevant).is_err()
4325        );
4326
4327        assert!(validate_native_authority_assignment_targets(
4328            &["peer-1".to_string(), "peer-1".to_string()],
4329            &[],
4330        )
4331        .is_err());
4332        assert!(validate_native_authority_assignment_targets(
4333            &[],
4334            &["entity-1".to_string(), "entity-1".to_string()],
4335        )
4336        .is_err());
4337    }
4338
4339    #[tokio::test]
4340    async fn authority_composition_starts_admitted_peer_owner_without_hosted_work() {
4341        let provider = Arc::new(RecordingGrantProvider::default());
4342        let capabilities = NativeCapabilities::new(
4343            crate::test_constants::TEST_API_KEY,
4344            "service-device",
4345            "native",
4346            provider.clone(),
4347        )
4348        .unwrap();
4349        let handle = capabilities
4350            .join_room_with_architecture("authority-room", RoomArchitectureMode::Authority)
4351            .unwrap();
4352        let client = handle
4353            .compose_authority_client(
4354                crate::Client::builder(
4355                    crate::test_constants::TEST_API_KEY.to_string(),
4356                    Box::new(|| None),
4357                )
4358                .unwrap(),
4359            )
4360            .await
4361            .unwrap();
4362        assert_eq!(
4363            client.external_desired_peer_debug_snapshot().await["active"],
4364            true,
4365            "authority construction must initialize the existing Rust peer actor"
4366        );
4367        assert!(
4368            client.session_registry_active(),
4369            "admission must be closed before endpoint bind"
4370        );
4371        assert_eq!(
4372            handle
4373                .signaling
4374                .authority_client
4375                .lock()
4376                .await
4377                .as_ref()
4378                .unwrap()
4379                .scope,
4380            "v2:room:authority-room",
4381            "native and browser must gate the same capability scope"
4382        );
4383        assert!(
4384            provider.requests.lock().unwrap().is_empty(),
4385            "composition must be network-idle"
4386        );
4387        assert!(
4388            handle.signaling.authority_ticket_maintenance_delay().await > Duration::from_secs(60),
4389            "an unpublished authority must not spin a renewal timer"
4390        );
4391        let before = client.external_desired_peer_debug_snapshot().await;
4392        assert!(handle
4393            .compose_authority_client(
4394                crate::Client::builder(
4395                    crate::test_constants::TEST_API_KEY.to_string(),
4396                    Box::new(|| None)
4397                )
4398                .unwrap(),
4399            )
4400            .await
4401            .is_err());
4402        assert_eq!(
4403            client.external_desired_peer_debug_snapshot().await["generation"],
4404            before["generation"]
4405        );
4406        handle.close().await;
4407        assert_eq!(
4408            client.external_desired_peer_debug_snapshot().await["active"],
4409            false
4410        );
4411        assert!(
4412            client.session_registry_active(),
4413            "close must not enable tokenless ingress"
4414        );
4415        assert!(handle
4416            .compose_authority_client(
4417                crate::Client::builder(
4418                    crate::test_constants::TEST_API_KEY.to_string(),
4419                    Box::new(|| None)
4420                )
4421                .unwrap(),
4422            )
4423            .await
4424            .is_err());
4425    }
4426
4427    #[tokio::test]
4428    async fn disconnected_authority_expires_admission_without_network_retry() {
4429        authority_admission_expires_without_extra_network_work(false).await;
4430    }
4431
4432    #[tokio::test]
4433    async fn connecting_authority_expires_admission_without_restarting_grant() {
4434        authority_admission_expires_without_extra_network_work(true).await;
4435    }
4436
4437    async fn authority_admission_expires_without_extra_network_work(connecting: bool) {
4438        let blocked = Arc::new(tokio::sync::Notify::new());
4439        let provider = Arc::new(RecordingGrantProvider {
4440            blocked: connecting.then(|| blocked.clone()),
4441            ..Default::default()
4442        });
4443        let handle = NativeCapabilities::new(
4444            crate::test_constants::TEST_API_KEY, "service-device", "native", provider.clone(),
4445        ).unwrap().join_room_with_architecture("expiry-room", RoomArchitectureMode::Authority).unwrap();
4446        let client = handle.compose_authority_client(crate::Client::builder(
4447            crate::test_constants::TEST_API_KEY.to_string(), Box::new(|| None),
4448        ).unwrap()).await.unwrap();
4449        let adapter = &handle.signaling;
4450        let scope = "v2:room:expiry-room";
4451        let peer = iroh::SecretKey::from_bytes(&[37; 32]).public().to_string();
4452        let expiry = now_ms() + 150;
4453        client.session_token_registry.update_scope_peer_admission(
4454            scope, 1, now_ms() + 60_000, vec![peer.clone()],
4455        ).unwrap();
4456        client.register_token_until("old-token".into(), scope.into(), 1, expiry);
4457        client.session_token_registry.validate_and_consume_for_authenticated_peer(
4458            "old-token", Some("expired-connection"), None, Some(&peer),
4459        ).unwrap();
4460        client.connection_manager.upsert_pending(
4461            "expired-connection".into(), None, Some("peer".into()), None,
4462        ).await;
4463        client.connection_manager.set_connected("expired-connection", None).await;
4464        let fresh_ticket = crate::session_token::expiring_ticket(
4465            "local-endpoint-ticket", "fresh-token", scope, 1, Some(now_ms() + 900_000),
4466        );
4467        *adapter.shared.desired.write().await = Some(DesiredPresence {
4468            user_id: "expiry-room".into(), local_node_id: peer,
4469            ticket: fresh_ticket, device_name: "Service".into(), metadata: None,
4470            ttl_ms: 60_000, online: true,
4471        });
4472        adapter.authority_client.lock().await.as_mut().unwrap()
4473            .retiring_tickets.push(("old-token".into(), expiry));
4474        let (reply, _publication) = oneshot::channel();
4475        if connecting {
4476            let desired = adapter.shared.desired.read().await.clone().unwrap();
4477            adapter.commands.send(Command::Publish { desired, reply }).await.unwrap();
4478            tokio::time::timeout(Duration::from_secs(1), async {
4479                while provider.requests.lock().unwrap().is_empty() {
4480                    tokio::time::sleep(Duration::from_millis(5)).await;
4481                }
4482            }).await.expect("actor must enter the single blocked grant request");
4483        } else {
4484            // Wake the parked actor without spending a network attempt.
4485            assert!(adapter.request_unit(|reply| Command::Patch {
4486                patch: DevicePatch::default(), reply,
4487            }).await.is_err());
4488        }
4489        let retired = tokio::time::timeout(Duration::from_secs(1), async {
4490            loop {
4491                let state = client.connection_manager.get_by_connection_id("expired-connection").await;
4492                if state.is_some_and(|record| record.state == crate::connection_manager::ConnectionState::Closed) {
4493                    break;
4494                }
4495                tokio::time::sleep(Duration::from_millis(10)).await;
4496            }
4497        }).await;
4498        let remaining = adapter.authority_client.lock().await.as_ref().unwrap().retiring_tickets.len();
4499        let requests = provider.requests.lock().unwrap().len();
4500        blocked.notify_one();
4501        handle.close().await;
4502        assert!(retired.is_ok(), "disconnected actor deferred expired admission retirement until another connect");
4503        assert_eq!(remaining, 0);
4504        assert_eq!(requests, usize::from(connecting), "expiry must not spend an extra network attempt");
4505    }
4506
4507    #[tokio::test]
4508    async fn authority_routes_reach_rust_once_and_withdraw_without_roster_dialing() {
4509        let provider = Arc::new(RecordingGrantProvider::default());
4510        let handle = NativeCapabilities::new(
4511            crate::test_constants::TEST_API_KEY,
4512            "service-device",
4513            "native",
4514            provider.clone(),
4515        )
4516        .unwrap()
4517        .join_room_with_architecture("authority-room", RoomArchitectureMode::Authority)
4518        .unwrap();
4519        let client = handle
4520            .compose_authority_client(
4521                crate::Client::builder(
4522                    crate::test_constants::TEST_API_KEY.to_string(),
4523                    Box::new(|| None),
4524                )
4525                .unwrap(),
4526            )
4527            .await
4528            .unwrap();
4529        let adapter = &handle.signaling;
4530        let now = now_ms();
4531        let member = |id: &str| GatewayMember {
4532            user_id: None,
4533            device_id: id.into(),
4534            device_name: id.into(),
4535            platform_type: "native".into(),
4536            metadata: None,
4537            capabilities: None,
4538            online: true,
4539            updated_at_ms: now as i64,
4540            expires_at_ms: (now + 60_000) as i64,
4541        };
4542        let route = |id: &str, seed: u8, scope: &str| {
4543            let key = iroh::SecretKey::from_bytes(&[seed; 32]);
4544            let ticket =
4545                iroh_tickets::endpoint::EndpointTicket::new(iroh::EndpointAddr::new(key.public()))
4546                    .to_string();
4547            GatewayDevice {
4548                user_id: None,
4549                device_id: id.into(),
4550                runtime_instance_id: format!("runtime-{id}"),
4551                node_id: key.public().to_string(),
4552                device_name: id.into(),
4553                platform_type: "native".into(),
4554                ticket: crate::session_token::build_ticket(&ticket, "test-token", scope, 16),
4555                metadata: None,
4556                capabilities: None,
4557                excluded_peers: Vec::new(),
4558                online: true,
4559                updated_at_ms: now as i64,
4560                expires_at_ms: (now + 60_000) as i64,
4561            }
4562        };
4563        handle_frame(
4564            adapter,
4565            ServerFrame::MembershipSnapshot {
4566                members: vec![member("active"), member("backup"), member("roster-only")],
4567            },
4568        )
4569        .await
4570        .unwrap();
4571        assert_eq!(
4572            client.external_desired_peer_debug_snapshot().await["peers"],
4573            serde_json::json!([])
4574        );
4575        *adapter.shared.active_grant.write().await = Some(("grant-current".into(), now + 60_000));
4576        let lease = GatewayTopologyLease {
4577            schema_version: 2,
4578            topology_revision: 1,
4579            previous_topology_revision: None,
4580            avenue: adapter.avenue.clone(),
4581            grant_jti: "grant-current".into(),
4582            expires_at_ms: now + 30_000,
4583            architecture: GatewayArchitectureLease {
4584                policy_version: "room-architecture-v1".into(),
4585                epoch: 1,
4586                previous_architecture_epoch: None,
4587                requested_mode: RoomArchitectureMode::Authority,
4588                effective_mode: EffectiveRoomArchitecture::Authority,
4589                phase: RoomArchitecturePhase::Settled,
4590                reason: RoomArchitectureReason::Manual,
4591                held_credits_microusd: 25_000,
4592                quote_expires_at_ms: now + 30_000,
4593                reservation_id: Some("reservation:test".into()),
4594            },
4595            group_encryption: None,
4596            active: vec![route("active", 11, "v2:room:authority-room")],
4597            admission_peers: Some(vec![
4598                GatewayAdmissionPeer {
4599                    device_id: "active".into(),
4600                    node_id: route("active", 11, "v2:room:authority-room").node_id,
4601                },
4602                GatewayAdmissionPeer {
4603                    device_id: "backup".into(),
4604                    node_id: route("backup", 12, "v2:room:authority-room").node_id,
4605                },
4606            ]),
4607            backups: vec![route("backup", 12, "v2:room:authority-room")],
4608        };
4609        let mut wrong_scope = lease.clone();
4610        let mut missing_admission = lease.clone();
4611        missing_admission.admission_peers = None;
4612        assert!(accept_topology_lease(adapter, missing_admission)
4613            .await
4614            .is_err());
4615        wrong_scope.active = vec![route("active", 11, "user-device")];
4616        assert!(accept_topology_lease(adapter, wrong_scope).await.is_err());
4617        assert_eq!(*adapter.shared.topology_revision.read().await, 0);
4618        accept_topology_lease(adapter, lease.clone()).await.unwrap();
4619        let accepted = client.external_desired_peer_debug_snapshot().await;
4620        assert_eq!(accepted["peers"].as_array().unwrap().len(), 1);
4621        assert_eq!(accepted["peers"][0]["deviceId"], "active");
4622        assert_eq!(accepted["peers"][0]["sessionId"], "runtime-active");
4623        accept_topology_lease(adapter, lease.clone()).await.unwrap();
4624        handle_frame(adapter, ServerFrame::Pong).await.unwrap();
4625        assert_eq!(
4626            client.external_desired_peer_debug_snapshot().await["revision"],
4627            accepted["revision"]
4628        );
4629        handle_frame(
4630            adapter,
4631            ServerFrame::MembershipChanged {
4632                operation: "delete".into(),
4633                member: member("active"),
4634            },
4635        )
4636        .await
4637        .unwrap();
4638        assert_eq!(
4639            client.external_desired_peer_debug_snapshot().await["peers"],
4640            serde_json::json!([])
4641        );
4642        handle_frame(
4643            adapter,
4644            ServerFrame::MembershipChanged {
4645                operation: "upsert".into(),
4646                member: member("active"),
4647            },
4648        )
4649        .await
4650        .unwrap();
4651        assert_eq!(
4652            client.external_desired_peer_debug_snapshot().await["peers"],
4653            serde_json::json!([]),
4654            "rejoining the roster must not resurrect the previous route"
4655        );
4656        handle.close().await;
4657        let mut late = lease;
4658        late.topology_revision = 2;
4659        late.previous_topology_revision = Some(1);
4660        assert!(
4661            accept_topology_lease(adapter, late).await.is_err(),
4662            "a late lease cannot restore authority admission after close"
4663        );
4664        assert_eq!(
4665            client.external_desired_peer_debug_snapshot().await["active"],
4666            false
4667        );
4668        assert_eq!(
4669            client.external_desired_peer_debug_snapshot().await["peers"],
4670            serde_json::json!([])
4671        );
4672        assert!(provider.requests.lock().unwrap().is_empty());
4673    }
4674
4675    #[tokio::test]
4676    async fn authority_private_routes_exchange_protected_messages_on_loopback() {
4677        authority_loopback_payload(None, AuthorityLoopbackScenario::Renew).await;
4678    }
4679
4680    #[tokio::test]
4681    async fn authority_server_empty_routes_exchange_protected_messages_on_loopback() {
4682        // Gateway authority services receive no outbound assignment. Exercise
4683        // both endpoint orderings without teaching the service to dial rosters.
4684        for server_is_lower in [true, false] {
4685            authority_loopback_payload(Some(server_is_lower), AuthorityLoopbackScenario::Renew)
4686                .await;
4687        }
4688    }
4689
4690    #[tokio::test]
4691    async fn authority_fresh_bearer_reopens_retired_peer_without_server_dialing() {
4692        for server_is_lower in [true, false] {
4693            authority_loopback_payload(
4694                Some(server_is_lower),
4695                AuthorityLoopbackScenario::RetireAndRenew,
4696            )
4697            .await;
4698        }
4699    }
4700
4701    #[tokio::test]
4702    async fn authority_gateway_revocation_denies_held_stream_and_current_route() {
4703        for server_is_lower in [true, false] {
4704            authority_loopback_payload(
4705                Some(server_is_lower),
4706                AuthorityLoopbackScenario::GatewayRevoked,
4707            )
4708            .await;
4709        }
4710    }
4711
4712    enum AuthorityLoopbackScenario {
4713        Renew,
4714        RetireAndRenew,
4715        GatewayRevoked,
4716    }
4717
4718    async fn authority_loopback_payload(
4719        server_is_lower: Option<bool>,
4720        scenario: AuthorityLoopbackScenario,
4721    ) {
4722        let retire_before_renewal = matches!(scenario, AuthorityLoopbackScenario::RetireAndRenew);
4723        let provider = Arc::new(RecordingGrantProvider::default());
4724        let mut peers = Vec::new();
4725        for id in ["service-a", "service-b"] {
4726            let handle = NativeCapabilities::new(
4727                crate::test_constants::TEST_API_KEY,
4728                id,
4729                "native",
4730                provider.clone(),
4731            )
4732            .unwrap()
4733            .join_room_with_architecture("loopback-room", RoomArchitectureMode::Authority)
4734            .unwrap();
4735            let client = handle
4736                .compose_authority_client(
4737                    crate::Client::builder(
4738                        crate::test_constants::TEST_API_KEY.to_string(),
4739                        Box::new(|| None),
4740                    )
4741                    .unwrap()
4742                    .transport_config(crate::client::TransportConfig {
4743                        relay: false,
4744                        webrtc: None,
4745                        moq: None,
4746                        ..Default::default()
4747                    }),
4748                )
4749                .await
4750                .unwrap();
4751            // Minimal has no n0 lookup/publisher. Only these explicitly bound
4752            // loopback sockets are used; the fake grant provider is never called.
4753            let endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
4754                .relay_mode(iroh::RelayMode::Disabled)
4755                .alpns(vec![crate::native_node::PlutoniumProtocol::ALPN.to_vec()])
4756                .bind_addr("127.0.0.1:0".parse::<std::net::SocketAddr>().unwrap())
4757                .unwrap()
4758                .bind()
4759                .await
4760                .unwrap();
4761            client
4762                .adopt_endpoint_with_router_mode(endpoint.clone(), true)
4763                .await
4764                .unwrap();
4765            let ticket = client
4766                .endpoint_ticket_with_token("v2:room:loopback-room", 16)
4767                .await
4768                .unwrap();
4769            peers.push((handle, client, endpoint, ticket));
4770        }
4771        if let Some(server_is_lower) = server_is_lower {
4772            peers.sort_by_key(|peer| peer.2.id());
4773            if !server_is_lower {
4774                peers.reverse();
4775            }
4776        }
4777        let mut received_a = peers[0].1.subscribe_native_peer_data();
4778        let mut received_b = peers[1].1.subscribe_native_peer_data();
4779        let result = tokio::time::timeout(Duration::from_secs(15), async {
4780            let mut leases = Vec::new();
4781            for (local, remote) in [(0, 1), (1, 0)] {
4782                let adapter = &peers[local].0.signaling;
4783                let remote_adapter = &peers[remote].0.signaling;
4784                let now = now_ms();
4785                handle_frame(
4786                    adapter,
4787                    ServerFrame::MembershipSnapshot {
4788                        members: vec![GatewayMember {
4789                            user_id: None,
4790                            device_id: remote_adapter.device_id.clone(),
4791                            device_name: "Peer".into(),
4792                            platform_type: "native".into(),
4793                            metadata: None,
4794                            capabilities: None,
4795                            online: true,
4796                            updated_at_ms: now as i64,
4797                            expires_at_ms: (now + 60_000) as i64,
4798                        }],
4799                    },
4800                )
4801                .await?;
4802                *adapter.shared.active_grant.write().await =
4803                    Some(("local-test-grant".into(), now + 60_000));
4804                let mut lease = GatewayTopologyLease {
4805                    schema_version: 2,
4806                    topology_revision: 1,
4807                    previous_topology_revision: None,
4808                    avenue: adapter.avenue.clone(),
4809                    grant_jti: "local-test-grant".into(),
4810                    expires_at_ms: now + 30_000,
4811                    architecture: GatewayArchitectureLease {
4812                        policy_version: "room-architecture-v1".into(),
4813                        epoch: 1,
4814                        previous_architecture_epoch: None,
4815                        requested_mode: RoomArchitectureMode::Authority,
4816                        effective_mode: EffectiveRoomArchitecture::Authority,
4817                        phase: RoomArchitecturePhase::Settled,
4818                        reason: RoomArchitectureReason::Manual,
4819                        held_credits_microusd: 1,
4820                        quote_expires_at_ms: now + 30_000,
4821                        reservation_id: Some("reservation:local".into()),
4822                    },
4823                    group_encryption: None,
4824                    backups: Vec::new(),
4825                    admission_peers: Some(vec![GatewayAdmissionPeer {
4826                        device_id: remote_adapter.device_id.clone(),
4827                        node_id: peers[remote].2.id().to_string(),
4828                    }]),
4829                    active: vec![GatewayDevice {
4830                        user_id: None,
4831                        device_id: remote_adapter.device_id.clone(),
4832                        runtime_instance_id: remote_adapter.runtime_instance_id.clone(),
4833                        node_id: peers[remote].2.id().to_string(),
4834                        device_name: "Peer".into(),
4835                        platform_type: "native".into(),
4836                        ticket: peers[remote].3.clone(),
4837                        metadata: None,
4838                        capabilities: None,
4839                        excluded_peers: Vec::new(),
4840                        online: true,
4841                        updated_at_ms: now as i64,
4842                        expires_at_ms: (now + 60_000) as i64,
4843                    }],
4844                };
4845                if server_is_lower.is_some() && local == 0 {
4846                    lease.active.clear();
4847                }
4848                accept_topology_lease(adapter, lease.clone()).await?;
4849                leases.push(lease);
4850            }
4851            let (a, b) = loop {
4852                let a = peers[0]
4853                    .1
4854                    .connection_states()
4855                    .await
4856                    .into_iter()
4857                    .find(|state| state.routable);
4858                let b = peers[1]
4859                    .1
4860                    .connection_states()
4861                    .await
4862                    .into_iter()
4863                    .find(|state| state.routable);
4864                if let (Some(a), Some(b)) = (a, b) {
4865                    break (a, b);
4866                }
4867                tokio::time::sleep(Duration::from_millis(20)).await;
4868            };
4869            peers[0]
4870                .1
4871                .send_peer(&a.connection_id, b"authority-a-to-b")
4872                .await?;
4873            peers[1]
4874                .1
4875                .send_peer(&b.connection_id, b"authority-b-to-a")
4876                .await?;
4877            assert_eq!(received_a.recv().await?.payload, b"authority-b-to-a");
4878            assert_eq!(received_b.recv().await?.payload, b"authority-a-to-b");
4879            if matches!(scenario, AuthorityLoopbackScenario::GatewayRevoked) {
4880                let (_, _, mut held, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
4881                held.write_all(b"authorized-stream-prefix").await?;
4882                assert!(
4883                    peers[0]
4884                        .0
4885                        .signaling
4886                        .handle_gateway_failure(&gateway_connect_error(
4887                            "local network outage",
4888                            true
4889                        ),)
4890                        .await
4891                );
4892                held.write_all(b"still-authorized-during-outage").await?;
4893                assert!(
4894                    peers[0]
4895                        .1
4896                        .connection_state(&a.connection_id)
4897                        .await
4898                        .unwrap()
4899                        .routable
4900                );
4901                let error = handle_frame(
4902                    &peers[0].0.signaling,
4903                    ServerFrame::Error {
4904                        code: "credential-revoked".into(),
4905                        message: "local administrative revocation".into(),
4906                        retryable: false,
4907                        idempotency_key: None,
4908                        metadata: Default::default(),
4909                    },
4910                )
4911                .await
4912                .unwrap_err();
4913                assert!(!peers[0].0.signaling.handle_gateway_failure(&error).await);
4914                assert!(
4915                    held.write_all(b"revoked-held-stream").await.is_err(),
4916                    "gateway revocation must withdraw existing protected stream access"
4917                );
4918                assert!(peers[0]
4919                    .1
4920                    .send_peer(&a.connection_id, b"revoked-message")
4921                    .await
4922                    .is_err());
4923                assert!(peers[0]
4924                    .1
4925                    .connection_states()
4926                    .await
4927                    .iter()
4928                    .all(|state| !state.routable));
4929                assert!(
4930                    peers[0]
4931                        .0
4932                        .signaling
4933                        .authority_client
4934                        .lock()
4935                        .await
4936                        .as_ref()
4937                        .unwrap()
4938                        .closed
4939                );
4940                let mut late = leases[0].clone();
4941                late.previous_topology_revision = Some(late.topology_revision);
4942                late.topology_revision += 1;
4943                assert!(
4944                    accept_topology_lease(&peers[0].0.signaling, late)
4945                        .await
4946                        .is_err(),
4947                    "late gateway projection must not reopen the revoked room owner"
4948                );
4949                return anyhow::Ok(());
4950            }
4951            // An expiring room bearer must be replaceable without replacing
4952            // the transport or relying on the app to restart the service. In
4953            // the asymmetric case only the consumer has an outbound assignment;
4954            // renew the server ticket with both endpoint orderings.
4955            let presenter = usize::from(server_is_lower.is_some());
4956            let host = 1 - presenter;
4957            let host_connection_id = if host == 0 {
4958                &a.connection_id
4959            } else {
4960                &b.connection_id
4961            };
4962            let adapter = &peers[host].0.signaling;
4963            let initial_presence = DesiredPresence {
4964                user_id: "loopback-room".into(),
4965                local_node_id: peers[host].2.id().to_string(),
4966                ticket: peers[host].3.clone(),
4967                device_name: "service-b".into(),
4968                metadata: None,
4969                ttl_ms: 60_000,
4970                online: true,
4971            };
4972            *adapter.shared.desired.write().await = Some(initial_presence.clone());
4973            assert!(!adapter.maintain_authority_ticket().await?);
4974            assert!(adapter.authority_ticket_maintenance_delay().await > Duration::from_secs(60));
4975            let (ticket, suffix) = crate::session_token::split_ticket(&peers[host].3);
4976            let old = crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
4977            // Use a short-lived local fixture instead of sleeping 15 minutes.
4978            let short_expiry = now_ms() + 30_000;
4979            peers[host].1.register_token_until(
4980                old.token.clone(),
4981                old.scope.to_string(),
4982                old.max_connections,
4983                short_expiry,
4984            );
4985            adapter
4986                .shared
4987                .desired
4988                .write()
4989                .await
4990                .as_mut()
4991                .unwrap()
4992                .ticket = crate::session_token::expiring_ticket(
4993                ticket,
4994                &old.token,
4995                &old.scope,
4996                old.max_connections,
4997                Some(short_expiry),
4998            );
4999            assert!(adapter.authority_ticket_maintenance_delay().await <= Duration::from_millis(1));
5000            let mut pending = adapter.shared.desired.read().await.clone().unwrap();
5001            assert!(adapter.maintain_authority_publication(Some(&mut pending)).await?);
5002            assert_eq!(Some(&pending), adapter.shared.desired.read().await.as_ref(),
5003                "local renewal must preserve the pending publication acknowledgement");
5004            let replacement = adapter
5005                .shared
5006                .desired
5007                .read()
5008                .await
5009                .as_ref()
5010                .unwrap()
5011                .ticket
5012                .clone();
5013            assert!(
5014                !adapter.maintain_authority_ticket().await?,
5015                "unchanged ticket rotated twice"
5016            );
5017            assert_eq!(
5018                adapter
5019                    .authority_client
5020                    .lock()
5021                    .await
5022                    .as_ref()
5023                    .unwrap()
5024                    .retiring_tickets
5025                    .len(),
5026                1
5027            );
5028            let (ticket, suffix) = crate::session_token::split_ticket(&replacement);
5029            let replacement_payload =
5030                crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
5031            let replacement_fingerprint =
5032                crate::session_token::token_fingerprint(&replacement_payload.token);
5033            if retire_before_renewal {
5034                // An outage can outlive the old bearer. Retire that exact token
5035                // through its owner before publishing the already minted fresh
5036                // one; a logical tombstone must not prevent fresh Iroh proof.
5037                assert!(!peers[host].1.revoke_session_token(&old.token).is_empty());
5038                while peers[host]
5039                    .1
5040                    .connection_manager
5041                    .get_by_connection_id(host_connection_id)
5042                    .await
5043                    .is_some()
5044                {
5045                    tokio::time::sleep(Duration::from_millis(20)).await;
5046                }
5047            }
5048            leases[presenter].topology_revision += 1;
5049            leases[presenter].previous_topology_revision = Some(1);
5050            leases[presenter].active[0].ticket = replacement;
5051            accept_topology_lease(&peers[presenter].0.signaling, leases[presenter].clone()).await?;
5052            while peers[host]
5053                .1
5054                .session_token_registry
5055                .admission_fingerprint(host_connection_id)
5056                .as_deref()
5057                != Some(replacement_fingerprint.as_str())
5058            {
5059                tokio::time::sleep(Duration::from_millis(20)).await;
5060            }
5061            // Drive the overlap deadline directly; cleanup must retire only
5062            // the old bearer, not the newly admitted physical generation.
5063            adapter
5064                .authority_client
5065                .lock()
5066                .await
5067                .as_mut()
5068                .unwrap()
5069                .retiring_tickets[0]
5070                .1 = 0;
5071            assert!(!adapter.maintain_authority_ticket().await?);
5072            assert!(adapter
5073                .authority_client
5074                .lock()
5075                .await
5076                .as_ref()
5077                .unwrap()
5078                .retiring_tickets
5079                .is_empty());
5080            assert!(peers[host]
5081                .1
5082                .session_token_registry
5083                .validate_and_consume(&old.token)
5084                .is_err());
5085            // Admission is recorded before route settlement. Observe both
5086            // owners' readiness before inspecting the replacement generation.
5087            loop {
5088                let a_ready = peers[0].1.connection_state(&a.connection_id).await
5089                    .is_some_and(|state| state.routable);
5090                let b_ready = peers[1].1.connection_state(&b.connection_id).await
5091                    .is_some_and(|state| state.routable);
5092                if a_ready && b_ready { break; }
5093                tokio::time::sleep(Duration::from_millis(20)).await;
5094            }
5095            assert_eq!(
5096                peers[0]
5097                    .1
5098                    .connection_state(&a.connection_id)
5099                    .await
5100                    .unwrap()
5101                    .active_transport_stable_id
5102                    == a.active_transport_stable_id,
5103                !retire_before_renewal,
5104                "revoked physical generation must be replaced; healthy renewal must retain it"
5105            );
5106            assert_eq!(
5107                peers[1]
5108                    .1
5109                    .connection_state(&b.connection_id)
5110                    .await
5111                    .unwrap()
5112                    .active_transport_stable_id
5113                    == b.active_transport_stable_id,
5114                !retire_before_renewal
5115            );
5116            peers[0]
5117                .1
5118                .send_peer(&a.connection_id, b"after-room-token-rotation")
5119                .await?;
5120            assert_eq!(
5121                received_b.recv().await?.payload,
5122                b"after-room-token-rotation"
5123            );
5124            peers[1]
5125                .1
5126                .send_peer(&b.connection_id, b"reverse-after-room-token-rotation")
5127                .await?;
5128            assert_eq!(
5129                received_a.recv().await?.payload,
5130                b"reverse-after-room-token-rotation"
5131            );
5132            if server_is_lower.is_some() {
5133                assert!(
5134                    leases[0].active.is_empty(),
5135                    "service must not dial its clients"
5136                );
5137                for (index, connection_id) in [(0, &a.connection_id), (1, &b.connection_id)] {
5138                    let state = peers[index]
5139                        .1
5140                        .connection_state(connection_id)
5141                        .await
5142                        .unwrap();
5143                    assert!(state.routable);
5144                    let scopes = peers[index]
5145                        .1
5146                        .connection_manager
5147                        .get_scopes(connection_id)
5148                        .await;
5149                    assert!(!scopes.iter().any(|scope| scope == "user-device"));
5150                }
5151                return anyhow::Ok(());
5152            }
5153
5154            // The SDK must not steal a product's stream or lose ciphertext
5155            // while looking for its own channel, even with fragmented headers.
5156            let unclaimed = peers[1].1.incoming_streams().await?;
5157            assert!(
5158                unclaimed.try_recv().is_err(),
5159                "default messages escaped to the host"
5160            );
5161            let (_, _, mut send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
5162            let mut expected = crate::stream_metadata::encode_envelope("app/custom", None)?;
5163            expected.extend_from_slice(b"unclaimed application bytes");
5164            send.write_all(&expected[..2]).await?;
5165            send.write_all(&expected[2..7]).await?;
5166            send.write_all(&expected[7..]).await?;
5167            send.finish_and_wait_for_peer(Duration::from_secs(2))
5168                .await?;
5169            let incoming = unclaimed.recv().await?;
5170            let crate::native_node::IncomingStreamType::Bi(send, recv) = incoming.stream else {
5171                panic!("expected bidirectional application stream");
5172            };
5173            let (_, mut recv) = peers[1]
5174                .1
5175                .wrap_bi_with_prefix(
5176                    &incoming.endpoint_id,
5177                    incoming.transport_stable_id,
5178                    send,
5179                    recv,
5180                    &incoming.recv_prefix,
5181                )
5182                .await?;
5183            let mut actual = Vec::new();
5184            let mut buffer = [0; 32];
5185            loop {
5186                let read = recv.read(&mut buffer).await?;
5187                if read == 0 {
5188                    break;
5189                }
5190                actual.extend_from_slice(&buffer[..read]);
5191            }
5192            assert_eq!(actual, expected);
5193
5194            // Reject an oversized default frame before allocation, and keep
5195            // the logical connection usable for the next valid message.
5196            let (_, _, mut send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
5197            send.write_all(&crate::stream_metadata::encode_envelope(
5198                crate::stream_metadata::DEFAULT_PEER_CHANNEL_ID,
5199                None,
5200            )?)
5201            .await?;
5202            send.write_all(&u32::MAX.to_be_bytes()).await?;
5203            let _ = send.finish_and_wait_for_peer(Duration::from_secs(2)).await;
5204            peers[0]
5205                .1
5206                .send_peer(&a.connection_id, b"after-invalid-message")
5207                .await?;
5208            assert_eq!(received_b.recv().await?.payload, b"after-invalid-message");
5209            assert!(received_b.try_recv().is_err());
5210            assert!(
5211                unclaimed.try_recv().is_err(),
5212                "invalid reserved channel escaped to host"
5213            );
5214
5215            let protected = peers[0]
5216                .1
5217                .protect_outbound_application_payload(&a.connection_id, b"stale")?;
5218            let mut frame = vec![0];
5219            frame.extend_from_slice(&protected);
5220            assert!(peers[1]
5221                .1
5222                .open_inbound_frame(
5223                    &b.connection_id,
5224                    b.remote_node_id.as_deref(),
5225                    "iroh",
5226                    Some(u64::MAX),
5227                    &frame,
5228                )
5229                .await
5230                .is_err());
5231            assert!(received_b.try_recv().is_err());
5232            // A gateway outage does not deliver an expiry/cleanup event. Even
5233            // with the physical route and application key still installed,
5234            // the shared admission owner must deny an expired service bearer.
5235            let (_, _, mut live_send, _) = peers[0].1.open_peer_bi(&a.connection_id, None).await?;
5236            let envelope = crate::stream_metadata::encode_envelope("app/expiry", None)?;
5237            live_send.write_all(&envelope).await?;
5238            let incoming = unclaimed.recv().await?;
5239            let crate::native_node::IncomingStreamType::Bi(send, recv) = incoming.stream else {
5240                panic!("expected bidirectional application stream");
5241            };
5242            let (mut expired_send, mut expired_recv) = peers[1]
5243                .1
5244                .wrap_bi_with_prefix(
5245                    &incoming.endpoint_id,
5246                    incoming.transport_stable_id,
5247                    send,
5248                    recv,
5249                    &incoming.recv_prefix,
5250                )
5251                .await?;
5252            let mut header = vec![0; envelope.len()];
5253            assert_eq!(expired_recv.read(&mut header).await?, envelope.len());
5254            assert_eq!(header, envelope);
5255            // The read begins while authorized and resumes after expiry.
5256            let mut denied_bytes = [0x55; 32];
5257            {
5258                let pending_read = expired_recv.read(&mut denied_bytes);
5259                tokio::pin!(pending_read);
5260                assert!(
5261                    tokio::time::timeout(Duration::from_millis(20), &mut pending_read)
5262                        .await
5263                        .is_err()
5264                );
5265                peers[1].1.register_token_until(
5266                    replacement_payload.token.clone(),
5267                    replacement_payload.scope.to_string(),
5268                    replacement_payload.max_connections,
5269                    0,
5270                );
5271                live_send.write_all(b"must-not-be-delivered").await?;
5272                assert_eq!(
5273                    pending_read
5274                        .await
5275                        .expect_err("open stream must recheck admission after awaiting data")
5276                        .kind(),
5277                    std::io::ErrorKind::PermissionDenied
5278                );
5279            }
5280            assert_eq!(
5281                denied_bytes, [0x55; 32],
5282                "denied reads cannot modify the consumer buffer"
5283            );
5284            let mut expired_recv =
5285                expired_recv.with_plaintext_prefix(b"buffered private prefix".iter().copied());
5286            assert_eq!(
5287                expired_recv
5288                    .read(&mut denied_bytes)
5289                    .await
5290                    .unwrap_err()
5291                    .kind(),
5292                std::io::ErrorKind::PermissionDenied
5293            );
5294            assert_eq!(
5295                expired_send
5296                    .write_all(b"must-not-be-sent")
5297                    .await
5298                    .expect_err("open stream must recheck admission before sending")
5299                    .kind(),
5300                std::io::ErrorKind::PermissionDenied
5301            );
5302            assert!(peers[1]
5303                .1
5304                .ensure_native_stream_admitted(&b.connection_id, b.remote_node_id.as_deref(), None,)
5305                .is_err());
5306            let error = peers[1]
5307                .1
5308                .open_inbound_frame(
5309                    &b.connection_id,
5310                    b.remote_node_id.as_deref(),
5311                    "iroh",
5312                    b.active_transport_stable_id,
5313                    &frame,
5314                )
5315                .await
5316                .expect_err("expiry blocks protected messages without a gateway event");
5317            assert!(error.to_string().contains("not admitted"), "{error:#}");
5318            assert!(received_b.try_recv().is_err());
5319            let recovered_ticket = peers[1]
5320                .1
5321                .endpoint_ticket_with_token("v2:room:loopback-room", 16)
5322                .await?;
5323            let (ticket, suffix) = crate::session_token::split_ticket(&recovered_ticket);
5324            let recovered = crate::session_token::decode_payload(ticket, suffix.unwrap()).unwrap();
5325            leases[0].previous_topology_revision = Some(leases[0].topology_revision);
5326            leases[0].topology_revision += 1;
5327            leases[0].active[0].ticket = recovered_ticket;
5328            accept_topology_lease(&peers[0].0.signaling, leases[0].clone()).await?;
5329            while peers[1]
5330                .1
5331                .session_token_registry
5332                .admission_fingerprint(&b.connection_id)
5333                .as_deref()
5334                != Some(crate::session_token::token_fingerprint(&recovered.token).as_str())
5335            {
5336                tokio::time::sleep(Duration::from_millis(20)).await;
5337            }
5338            peers[0]
5339                .1
5340                .send_peer(&a.connection_id, b"after-expired-admission-recovery")
5341                .await?;
5342            assert_eq!(
5343                received_b.recv().await?.payload,
5344                b"after-expired-admission-recovery"
5345            );
5346            assert_eq!(
5347                expired_recv
5348                    .read(&mut denied_bytes)
5349                    .await
5350                    .unwrap_err()
5351                    .kind(),
5352                std::io::ErrorKind::PermissionDenied,
5353                "reauthorization cannot revive a stream that observed denial"
5354            );
5355            anyhow::Ok(())
5356        })
5357        .await;
5358        let states_a = peers[0].1.connection_states().await;
5359        let states_b = peers[1].1.connection_states().await;
5360        for (handle, _, endpoint, _) in &peers {
5361            handle.close().await;
5362            endpoint.close().await;
5363        }
5364        assert!(provider.requests.lock().unwrap().is_empty());
5365        result
5366            .unwrap_or_else(|_| {
5367                panic!("protected authority route timeout: server_is_lower={server_is_lower:?} a={states_a:?} b={states_b:?}")
5368            })
5369            .unwrap();
5370    }
5371
5372    #[tokio::test]
5373    async fn native_capability_handle_is_inactive_until_runtime_presence_and_closes_once() {
5374        let provider = Arc::new(RecordingGrantProvider::default());
5375        let capabilities = NativeCapabilities::new(
5376            crate::test_constants::TEST_API_KEY,
5377            "device-native",
5378            "native",
5379            provider.clone(),
5380        )
5381        .expect("capability namespace")
5382        .with_testing_endpoint("https://gateway.example.test")
5383        .expect("testing endpoint");
5384        let handle = capabilities
5385            .join_room_with_architecture("room-1", RoomArchitectureMode::Auto)
5386            .expect("capability handle");
5387
5388        assert_eq!(handle.kind(), NativeCapabilityKind::Room);
5389        assert_eq!(handle.id(), "room-1");
5390        assert!(!handle.is_closed());
5391        assert!(provider.requests.lock().expect("requests").is_empty());
5392
5393        handle.close().await;
5394        handle.close().await;
5395        assert!(handle.is_closed());
5396        assert!(provider.requests.lock().expect("requests").is_empty());
5397    }
5398
5399    #[tokio::test]
5400    async fn room_architecture_intent_reaches_the_native_grant_owner() {
5401        let provider = Arc::new(RecordingGrantProvider::default());
5402        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
5403            endpoint: "https://gateway.example.test".to_string(),
5404            app_tag: "app_native_test".to_string(),
5405            device_id: "device-1".to_string(),
5406            platform_type: "desktop".to_string(),
5407            avenue: NativeCoordinationAvenue {
5408                kind: "room".to_string(),
5409                id: "room-1".to_string(),
5410            },
5411            architecture: Some(RoomArchitectureMode::Managed),
5412            room_delivery: Some(RoomDelivery::Reliable),
5413            grant_provider: provider.clone(),
5414            #[cfg(feature = "managed-group-encryption")]
5415            managed_group_signer: None,
5416            #[cfg(feature = "managed-group-encryption")]
5417            managed_group_store: None,
5418        })
5419        .expect("adapter");
5420        let desired = DesiredPresence {
5421            user_id: "room-1".to_string(),
5422            local_node_id: "node-1".to_string(),
5423            ticket: "ticket-1".to_string(),
5424            device_name: "Native".to_string(),
5425            metadata: None,
5426            ttl_ms: 900_000,
5427            online: true,
5428        };
5429
5430        adapter
5431            .mint_credential(&desired, "presence", None)
5432            .await
5433            .expect("grant");
5434
5435        let requests = provider.requests.lock().expect("requests");
5436        assert_eq!(requests.len(), 1);
5437        assert_eq!(
5438            requests[0].architecture,
5439            Some(RoomArchitectureMode::Managed)
5440        );
5441        assert_eq!(requests[0].room_delivery, Some(RoomDelivery::Reliable));
5442    }
5443
5444    #[tokio::test]
5445    async fn native_authority_assignment_is_local_monotonic_and_bounded() {
5446        let provider = Arc::new(RecordingGrantProvider::default());
5447        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
5448            endpoint: "https://gateway.example.test".to_string(),
5449            app_tag: "app_native_test".to_string(),
5450            device_id: "device-local".to_string(),
5451            platform_type: "desktop".to_string(),
5452            avenue: NativeCoordinationAvenue {
5453                kind: "room".to_string(),
5454                id: "room-authority".to_string(),
5455            },
5456            architecture: Some(RoomArchitectureMode::Authority),
5457            room_delivery: Some(RoomDelivery::Reliable),
5458            grant_provider: provider,
5459            #[cfg(feature = "managed-group-encryption")]
5460            managed_group_signer: None,
5461            #[cfg(feature = "managed-group-encryption")]
5462            managed_group_store: None,
5463        })
5464        .expect("adapter");
5465        adapter
5466            .shared
5467            .room_architecture
5468            .send_replace(Some(RoomArchitectureSnapshot {
5469                requested: RoomArchitectureMode::Authority,
5470                effective: EffectiveRoomArchitecture::Authority,
5471                epoch: 2,
5472                phase: RoomArchitecturePhase::Settled,
5473                reason: RoomArchitectureReason::Manual,
5474                held_credits_usd: 0.01,
5475                quote_expires_at_ms: now_ms() + 30_000,
5476            }));
5477        let assignment = NativeAuthorityInterestAssignment {
5478            policy_version: "room-authority-assignment-v1".to_string(),
5479            service_id: "service-primary".to_string(),
5480            generation: 2,
5481            revision: 1,
5482            subject_device_id: "device-local".to_string(),
5483            shard_id: "zone-a".to_string(),
5484            priority_device_ids: vec!["device-nearby".to_string()],
5485            relevant_entity_ids: vec!["entity-player-1".to_string()],
5486            issued_at_ms: now_ms(),
5487            expires_at_ms: now_ms() + 60_000,
5488        };
5489        let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([0_u8; 64]);
5490
5491        for _ in 0..2 {
5492            handle_frame(
5493                &adapter,
5494                ServerFrame::AuthorityAssignment {
5495                    assignment: assignment.clone(),
5496                    signature: signature.clone(),
5497                },
5498            )
5499            .await
5500            .expect("current assignment");
5501        }
5502        assert_eq!(
5503            adapter.authority_assignment().unwrap().assignment,
5504            assignment
5505        );
5506
5507        let stale = NativeAuthorityInterestAssignment {
5508            generation: 1,
5509            ..assignment.clone()
5510        };
5511        assert!(handle_frame(
5512            &adapter,
5513            ServerFrame::AuthorityAssignment {
5514                assignment: stale,
5515                signature: signature.clone(),
5516            },
5517        )
5518        .await
5519        .is_err());
5520        let wrong_subject = NativeAuthorityInterestAssignment {
5521            revision: 2,
5522            subject_device_id: "device-other".to_string(),
5523            ..assignment
5524        };
5525        assert!(handle_frame(
5526            &adapter,
5527            ServerFrame::AuthorityAssignment {
5528                assignment: wrong_subject,
5529                signature,
5530            },
5531        )
5532        .await
5533        .is_err());
5534    }
5535
5536    #[test]
5537    fn reconnect_is_bounded_and_capped() {
5538        assert_eq!(reconnect_delay(1), Duration::from_millis(250));
5539        assert_eq!(reconnect_delay(12), MAX_RECONNECT_DELAY);
5540        assert_eq!(MAX_RECONNECT_ATTEMPTS, 12);
5541    }
5542
5543    #[test]
5544    fn authority_return_retry_never_resets_budget_or_retries_terminal_errors() {
5545        let missing =
5546            gateway_server_error("room-authority-unavailable".into(), "local".into(), true);
5547        let mut available = false;
5548        assert_eq!(
5549            gateway_reconnect_delay(10, &missing, &mut available),
5550            Some(MAX_RECONNECT_DELAY),
5551            "initial publication cannot accelerate a host that has never been available"
5552        );
5553        available = true;
5554        for attempt in 1..=13 {
5555            let error = if attempt < 10 {
5556                gateway_connect_error("network unavailable", true)
5557            } else {
5558                gateway_server_error("room-authority-unavailable".into(), "local".into(), true)
5559            };
5560            let expected = match attempt {
5561                10 => Some(reconnect_delay(1)),
5562                12.. => None,
5563                _ => Some(reconnect_delay(attempt)),
5564            };
5565            assert_eq!(
5566                gateway_reconnect_delay(attempt, &error, &mut available),
5567                expected
5568            );
5569        }
5570        assert!(!available);
5571        for (code, retryable) in [
5572            ("credential-revoked", false),
5573            ("room-architecture-conflict", false),
5574            ("room-authority-unavailable", false),
5575            ("budget-renewal-required", true),
5576        ] {
5577            available = true;
5578            let error = gateway_server_error(code.into(), "local".into(), retryable);
5579            assert_eq!(
5580                gateway_reconnect_delay(10, &error, &mut available),
5581                retryable.then_some(MAX_RECONNECT_DELAY)
5582            );
5583            assert!(
5584                available,
5585                "unrelated/permanent failures must not consume the return allowance"
5586            );
5587        }
5588        available = true;
5589        assert_eq!(gateway_reconnect_delay(12, &missing, &mut available), None);
5590        assert!(
5591            available,
5592            "even an unused allowance cannot bypass the attempt cap"
5593        );
5594        let untyped = gateway_connect_error("room-authority-unavailable", true);
5595        assert_eq!(
5596            gateway_reconnect_delay(10, &untyped, &mut available),
5597            Some(MAX_RECONNECT_DELAY)
5598        );
5599    }
5600
5601    #[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
5602    #[tokio::test]
5603    async fn recovered_gateway_retries_host_return_without_another_capped_delay() {
5604        gateway_close_actor_scenario(false).await;
5605    }
5606
5607    #[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
5608    #[tokio::test]
5609    async fn revoked_gateway_socket_stops_authority_without_minting_another_grant() {
5610        gateway_close_actor_scenario(true).await;
5611    }
5612
5613    #[cfg(all(feature = "testing-endpoints", feature = "test-relay-client"))]
5614    async fn gateway_close_actor_scenario(revoked: bool) {
5615        // Real actor and loopback authentication; virtual time skips only the
5616        // failed grant phase. The relay-test feature supplies Tokio test time.
5617        struct OutageProvider {
5618            inner: RecordingGrantProvider,
5619            failures: std::sync::atomic::AtomicU8,
5620            calls: std::sync::atomic::AtomicU8,
5621        }
5622        #[async_trait]
5623        impl NativeGatewayGrantProvider for OutageProvider {
5624            async fn grant(
5625                &self,
5626                request: NativeGatewayGrantRequest,
5627            ) -> Result<NativeGatewayGrant> {
5628                self.calls.fetch_add(1, Ordering::SeqCst);
5629                if self
5630                    .failures
5631                    .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
5632                        remaining.checked_sub(1)
5633                    })
5634                    .is_ok()
5635                {
5636                    bail!("local injected control-plane outage");
5637                }
5638                self.inner.grant(request).await
5639            }
5640        }
5641        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
5642        let endpoint = format!("http://{}", listener.local_addr().unwrap());
5643        let provider = Arc::new(OutageProvider {
5644            inner: RecordingGrantProvider {
5645                gateway_url: Some(endpoint.clone()),
5646                token: Some(format!(
5647                    "local.{}.test",
5648                    base64::engine::general_purpose::URL_SAFE_NO_PAD
5649                        .encode(br#"{"jti":"local-outage-grant"}"#)
5650                )),
5651                ..Default::default()
5652            },
5653            failures: 0.into(),
5654            calls: 0.into(),
5655        });
5656        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
5657            endpoint,
5658            app_tag: "app_native_test".into(),
5659            device_id: "service-consumer".into(),
5660            platform_type: "native".into(),
5661            avenue: NativeCoordinationAvenue {
5662                kind: "room".into(),
5663                id: "local-room".into(),
5664            },
5665            architecture: Some(RoomArchitectureMode::Authority),
5666            room_delivery: None,
5667            grant_provider: provider.clone(),
5668            #[cfg(feature = "managed-group-encryption")]
5669            managed_group_signer: None,
5670            #[cfg(feature = "managed-group-encryption")]
5671            managed_group_store: None,
5672        })
5673        .unwrap();
5674        let (cut, cut_rx) = oneshot::channel();
5675        let (host_missing, host_missing_rx) = oneshot::channel();
5676        let server = tokio::spawn(async move {
5677            let mut cut_rx = Some(cut_rx);
5678            let mut host_missing = Some(host_missing);
5679            for connection in 0..3 {
5680                let (tcp, _) = listener.accept().await.unwrap();
5681                let mut socket = tokio_tungstenite::accept_hdr_async(tcp,
5682                    |_: &tokio_tungstenite::tungstenite::handshake::server::Request,
5683                     mut response: tokio_tungstenite::tungstenite::handshake::server::Response| {
5684                        response.headers_mut().insert("Sec-WebSocket-Protocol",
5685                            HeaderValue::from_static(GATEWAY_PROTOCOL));
5686                        Ok(response)
5687                    }).await.unwrap();
5688                let auth: serde_json::Value =
5689                    serde_json::from_str(socket.next().await.unwrap().unwrap().to_text().unwrap())
5690                        .unwrap();
5691                assert_eq!(auth["type"], "auth");
5692                let response = if connection == 1 {
5693                    serde_json::json!({"type":"error", "code":"room-authority-unavailable",
5694                        "message":"Host is returning", "retryable":true})
5695                } else {
5696                    serde_json::json!({"type":"ready", "budgetRemainingMicrousd":1000})
5697                };
5698                socket
5699                    .send(Message::Text(response.to_string().into()))
5700                    .await
5701                    .unwrap();
5702                if connection == 0 {
5703                    cut_rx.take().unwrap().await.unwrap();
5704                    let frame =
5705                        revoked.then(|| tokio_tungstenite::tungstenite::protocol::CloseFrame {
5706                            code: 4403.into(),
5707                            reason: "credential revoked".into(),
5708                        });
5709                    socket.close(frame).await.unwrap();
5710                    if revoked {
5711                        return;
5712                    }
5713                } else if connection == 1 {
5714                    host_missing.take().unwrap().send(()).unwrap();
5715                    socket.close(None).await.unwrap();
5716                } else {
5717                    // Keep the recovered fixture alive through the assertion.
5718                    // The runtime can immediately send presence or keepalive;
5719                    // receiving that frame is not a request to drop its socket.
5720                    while let Some(Ok(frame)) = socket.next().await {
5721                        match frame {
5722                            Message::Text(text) if text == "ping" => {
5723                                socket.send(Message::Text("pong".into())).await.unwrap();
5724                            }
5725                            Message::Ping(payload) => {
5726                                socket.send(Message::Pong(payload)).await.unwrap();
5727                            }
5728                            Message::Close(_) => break,
5729                            _ => {}
5730                        }
5731                    }
5732                }
5733            }
5734        });
5735        tokio::time::timeout(
5736            Duration::from_secs(3),
5737            adapter.update_presence(
5738                "local-room",
5739                &iroh::SecretKey::from_bytes(&[17; 32]).public().to_string(),
5740                "local-ticket",
5741                true,
5742                "Consumer",
5743                900_000,
5744                None,
5745            ),
5746        )
5747        .await
5748        .unwrap()
5749        .unwrap();
5750        if revoked {
5751            // This marker asserts the real socket actor invokes the existing
5752            // room owner. The separate Iroh loopback test proves payload denial
5753            // with a live composed client, permission lease and held stream.
5754            *adapter.authority_client.lock().await = Some(AuthorityPeerBinding {
5755                client: std::sync::Weak::new(),
5756                scope: "v2:room:local-room".into(),
5757                revision: 1,
5758                admission_revision: 1,
5759                last_payload: "[]".into(),
5760                retiring_tickets: Vec::new(),
5761                closed: false,
5762            });
5763            cut.send(()).unwrap();
5764            let stopped = tokio::time::timeout(Duration::from_secs(3), async {
5765                while !adapter
5766                    .authority_client
5767                    .lock()
5768                    .await
5769                    .as_ref()
5770                    .unwrap()
5771                    .closed
5772                {
5773                    tokio::time::sleep(Duration::from_millis(10)).await;
5774                }
5775            })
5776            .await;
5777            adapter.stop().await;
5778            server.abort();
5779            let _ = server.await;
5780            assert!(
5781                stopped.is_ok(),
5782                "revocation close did not reach authority owner"
5783            );
5784            assert_eq!(
5785                provider.calls.load(Ordering::SeqCst),
5786                1,
5787                "revoked socket must not mint a replacement grant"
5788            );
5789            return;
5790        }
5791        provider.failures.store(9, Ordering::SeqCst);
5792        cut.send(()).unwrap();
5793        tokio::time::timeout(Duration::from_secs(3), async {
5794            while provider.calls.load(Ordering::SeqCst) < 2 {
5795                tokio::time::sleep(Duration::from_millis(1)).await;
5796            }
5797        })
5798        .await
5799        .unwrap();
5800        tokio::time::pause();
5801        for attempt in 2..=9 {
5802            tokio::time::advance(reconnect_delay(attempt - 1) + Duration::from_millis(2)).await;
5803            for _ in 0..100 {
5804                tokio::task::yield_now().await;
5805            }
5806            assert_eq!(provider.calls.load(Ordering::SeqCst), attempt + 1);
5807        }
5808        tokio::time::advance(reconnect_delay(9) - Duration::from_millis(1)).await;
5809        tokio::time::resume();
5810        tokio::time::timeout(Duration::from_secs(3), host_missing_rx)
5811            .await
5812            .unwrap()
5813            .unwrap();
5814        let recovered = tokio::time::timeout(Duration::from_secs(3), async {
5815            while adapter.shared.applied_presence.read().await.is_none() {
5816                tokio::time::sleep(Duration::from_millis(10)).await;
5817            }
5818        })
5819        .await;
5820        adapter.stop().await;
5821        server.abort();
5822        let _ = server.await;
5823        assert!(
5824            recovered.is_ok(),
5825            "host return must not inherit another 30-second network backoff"
5826        );
5827        assert_eq!(
5828            provider.calls.load(Ordering::SeqCst),
5829            12,
5830            "initial connection plus eleven recovery attempts; never reset the outage budget"
5831        );
5832    }
5833
5834    #[test]
5835    fn endpoints_reject_query_credentials() {
5836        assert!(validate_endpoint(
5837            "endpoint",
5838            "https://gateway.example.test?token=secret".to_string(),
5839            true,
5840        )
5841        .is_err());
5842        assert!(validate_endpoint(
5843            "endpoint",
5844            "https://user:secret@gateway.example.test".to_string(),
5845            true,
5846        )
5847        .is_err());
5848        assert!(
5849            validate_endpoint("endpoint", "http://gateway.example.test".to_string(), true,)
5850                .is_err()
5851        );
5852        assert!(validate_endpoint("endpoint", "http://127.0.0.1:8787".to_string(), true,).is_ok());
5853        assert_eq!(normalized_origin_scheme("wss"), "https");
5854        assert_eq!(normalized_origin_scheme("ws"), "http");
5855    }
5856
5857    #[tokio::test]
5858    async fn grant_provider_owns_native_issue_and_in_place_refresh() {
5859        let provider = Arc::new(RecordingGrantProvider::default());
5860        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
5861            endpoint: "https://gateway.example.test".to_string(),
5862            app_tag: "app_native_test".to_string(),
5863            device_id: "device-1".to_string(),
5864            platform_type: "desktop".to_string(),
5865            avenue: NativeCoordinationAvenue {
5866                kind: "user".to_string(),
5867                id: "principal-1".to_string(),
5868            },
5869            architecture: None,
5870            room_delivery: None,
5871            grant_provider: provider.clone(),
5872            #[cfg(feature = "managed-group-encryption")]
5873            managed_group_signer: None,
5874            #[cfg(feature = "managed-group-encryption")]
5875            managed_group_store: None,
5876        })
5877        .expect("adapter");
5878        let desired = DesiredPresence {
5879            user_id: "principal-1".to_string(),
5880            local_node_id: "node-1".to_string(),
5881            ticket: "ticket-1".to_string(),
5882            device_name: "Native".to_string(),
5883            metadata: None,
5884            ttl_ms: 900_000,
5885            online: true,
5886        };
5887        let issued = adapter
5888            .mint_credential(&desired, "presence", None)
5889            .await
5890            .expect("issue");
5891        let refreshed = adapter
5892            .mint_credential(&desired, "presence", Some(issued.token.clone()))
5893            .await
5894            .expect("refresh");
5895        assert_eq!(issued.protocol_version, 2);
5896        assert_eq!(refreshed.protocol_version, 2);
5897        assert!(adapter
5898            .gateway_url(&issued)
5899            .expect("url")
5900            .contains("/v2/connect/route-1"));
5901        let requests = provider.requests.lock().expect("requests");
5902        assert_eq!(requests.len(), 2);
5903        assert_eq!(requests[0].refresh_grant, None);
5904        assert_eq!(requests[1].refresh_grant.as_deref(), Some("grant-2"));
5905        assert_eq!(requests[1].avenue.kind, "user");
5906        assert_eq!(requests[1].avenue.id, "principal-1");
5907    }
5908
5909    #[test]
5910    fn patch_merge_preserves_unrelated_fields() {
5911        let mut target = DevicePatch {
5912            device_name: Some("old".into()),
5913            capabilities: Some(DeviceCapabilities {
5914                can_host: true,
5915                can_sync: true,
5916                read_only: false,
5917            }),
5918            metadata: None,
5919            excluded_peers: None,
5920        };
5921        merge_patch(
5922            &mut target,
5923            DevicePatch {
5924                excluded_peers: Some(vec!["peer".into()]),
5925                ..DevicePatch::default()
5926            },
5927        );
5928        assert_eq!(target.device_name.as_deref(), Some("old"));
5929        assert_eq!(target.excluded_peers, Some(vec!["peer".into()]));
5930    }
5931
5932    #[test]
5933    fn permanent_gateway_errors_are_not_retried() {
5934        let permanent = gateway_connect_error("invalid payload", false);
5935        let transient = gateway_connect_error("provider unavailable", true);
5936        assert!(!is_retryable_gateway_error(&permanent));
5937        assert!(is_retryable_gateway_error(&transient));
5938    }
5939
5940    #[test]
5941    fn only_typed_gateway_revocation_is_terminal_authority_withdrawal() {
5942        let revoked = decode_frame(Message::Close(Some(
5943            tokio_tungstenite::tungstenite::protocol::CloseFrame {
5944                code: 4403.into(),
5945                reason: "credential revoked".into(),
5946            },
5947        )))
5948        .unwrap_err();
5949        assert!(is_gateway_revocation(&revoked));
5950        assert!(!is_retryable_gateway_error(&revoked));
5951        assert!(!is_gateway_revocation(
5952            &decode_frame(Message::Close(None)).unwrap_err()
5953        ));
5954        assert!(!is_gateway_revocation(&gateway_connect_error(
5955            "credential-revoked",
5956            false
5957        )));
5958        assert!(!is_gateway_revocation(&gateway_server_error(
5959            "budget-exhausted".into(),
5960            "local".into(),
5961            false
5962        )));
5963        let wire = gateway_server_error("credential-revoked".into(), "local".into(), true);
5964        assert!(
5965            !is_retryable_gateway_error(&wire),
5966            "revocation cannot grant a retry through a contradictory flag"
5967        );
5968        assert!(is_gateway_revocation(&wire.context("operation failed")));
5969        let certificate = anyhow::Error::new(ControlPlaneHttpError {
5970            service_error: None,
5971            status: 412,
5972            message: "device revoked".into(),
5973            reason: Some("device-certificate-revoked".into()),
5974        })
5975        .context("obtain native coordination grant");
5976        assert!(is_gateway_revocation(&certificate));
5977    }
5978
5979    #[tokio::test]
5980    async fn command_reply_preserves_revocation_for_the_socket_owner() {
5981        let (reply, result) = oneshot::channel::<Result<()>>();
5982        let error = complete_command(
5983            reply,
5984            Err(gateway_server_error(
5985                "credential-revoked".into(),
5986                "local".into(),
5987                false,
5988            )),
5989            "session.put",
5990        )
5991        .unwrap_err();
5992        let caller = result.await.unwrap().unwrap_err();
5993        assert!(is_gateway_revocation(&caller));
5994        assert!(!is_retryable_gateway_error(&caller));
5995        assert!(is_gateway_revocation(&error));
5996        assert!(!is_retryable_gateway_error(&error));
5997    }
5998
5999    #[test]
6000    fn only_typed_retryable_budget_exhaustion_renews_an_operation() {
6001        assert!(operation_requires_budget_renewal(
6002            "budget-renewal-required",
6003            true,
6004        ));
6005        assert!(!operation_requires_budget_renewal(
6006            "budget-renewal-required",
6007            false,
6008        ));
6009        assert!(!operation_requires_budget_renewal(
6010            "provider-budget-exhausted",
6011            true,
6012        ));
6013    }
6014
6015    #[tokio::test]
6016    async fn command_failure_reply_preserves_the_native_error_chain() {
6017        let (reply, response) = oneshot::channel();
6018        let returned = complete_command::<()>(
6019            reply,
6020            Err(anyhow!("control-plane detail").context("renewal failed")),
6021            "native device patch",
6022        )
6023        .expect_err("the actor must observe the operation failure");
6024        let caller = response
6025            .await
6026            .expect("reply")
6027            .expect_err("the caller must observe the operation failure");
6028
6029        assert_eq!(caller.to_string(), "renewal failed: control-plane detail");
6030        assert!(
6031            format!("{returned:#}").contains("renewal failed: control-plane detail"),
6032            "the actor should retain the same detailed chain"
6033        );
6034    }
6035
6036    #[tokio::test]
6037    async fn command_failure_reply_preserves_service_metadata_for_both_owners() {
6038        for code in crate::service_errors::SERVICE_ERROR_CODES {
6039            let metadata = serde_json::json!({
6040                "code": code, "retryable": false, "scope": "app",
6041                "operation": "presence.upsert", "requestId": "request-42",
6042                "retryAfterMs": 120000, "resetAt": 1800000000000_u64,
6043            });
6044            for http in [false, true] {
6045                let denial = if http {
6046                    anyhow::Error::new(ControlPlaneHttpError {
6047                        status: 429, message: "denied".into(), reason: None,
6048                        service_error: crate::service_errors::ServiceError::from_value(&metadata),
6049                    })
6050                } else {
6051                    gateway_server_error_metadata(code.to_string(), "denied".into(), false,
6052                        metadata.as_object().unwrap().clone())
6053                }.context("refresh failed");
6054                let (reply, response) = oneshot::channel();
6055                let actor = complete_command::<()>(reply, Err(denial), "presence.upsert").unwrap_err();
6056                let caller = response.await.unwrap().unwrap_err();
6057                for error in [&actor, &caller] {
6058                    let service = crate::service_errors::ServiceError::from_error(error).unwrap();
6059                    assert_eq!(service.code, *code);
6060                    assert_eq!(service.request_id.as_deref(), Some("request-42"));
6061                    assert_eq!(service.retry_after_ms, Some(120000));
6062                    assert_eq!(service.reset_at, Some(1800000000000));
6063                    assert!(!is_retryable_gateway_error(error));
6064                }
6065            }
6066        }
6067    }
6068
6069    #[tokio::test]
6070    async fn terminal_publication_preserves_service_denial_without_another_grant() {
6071        struct DeniedProvider(std::sync::atomic::AtomicUsize);
6072        #[async_trait]
6073        impl NativeGatewayGrantProvider for DeniedProvider {
6074            async fn grant(&self, _: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
6075                self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
6076                Err(anyhow::Error::new(ControlPlaneHttpError {
6077                    status: 429, message: "app cap reached".into(), reason: None,
6078                    service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
6079                        "code": "app-budget-exhausted", "retryable": false, "scope": "app",
6080                        "operation": "gateway.grant.issue", "requestId": "terminal-request",
6081                    })),
6082                }))
6083            }
6084        }
6085        let provider = Arc::new(DeniedProvider(0.into()));
6086        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6087            endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
6088            device_id: "device-1".into(), platform_type: "desktop".into(),
6089            avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
6090            architecture: None, room_delivery: None, grant_provider: provider.clone(),
6091            #[cfg(feature = "managed-group-encryption")]
6092            managed_group_signer: None,
6093            #[cfg(feature = "managed-group-encryption")]
6094            managed_group_store: None,
6095        }).unwrap();
6096        let desired = DesiredPresence {
6097            user_id: "principal-1".into(), local_node_id: "node-1".into(),
6098            ticket: "ticket-1".into(), device_name: "Native".into(),
6099            metadata: None, ttl_ms: 900000, online: true,
6100        };
6101        for _ in 0..3 {
6102            let (reply, response) = oneshot::channel();
6103            adapter.commands.send(Command::Publish { desired: desired.clone(), reply }).await.unwrap();
6104            let error = tokio::time::timeout(Duration::from_secs(2), response).await.unwrap().unwrap().unwrap_err();
6105            let service = crate::service_errors::ServiceError::from_error(&error).unwrap();
6106            assert_eq!(service.code, "app-budget-exhausted");
6107            assert_eq!(service.request_id.as_deref(), Some("terminal-request"));
6108            assert!(!is_retryable_gateway_error(&error));
6109        }
6110        assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
6111        adapter.stop().await;
6112    }
6113
6114    #[test]
6115    fn gateway_wire_service_errors_preserve_safe_metadata_and_cooldown() {
6116        for code in crate::service_errors::SERVICE_ERROR_CODES {
6117            let frame = decode_frame(Message::Text(serde_json::json!({
6118                "type": "error", "code": code, "message": "admission denied",
6119                "retryable": true, "retryAfterMs": 120000,
6120                "scope": "app", "operation": "gateway.grant.issue",
6121                "requestId": "request-1", "providerCost": 99,
6122                "documentationUrl": "https://untrusted.invalid"
6123            }).to_string().into())).unwrap().unwrap();
6124            let ServerFrame::Error { code: actual, message, retryable, metadata, .. } = frame else {
6125                panic!("expected service denial");
6126            };
6127            let error = gateway_server_error_metadata(actual, message, retryable, metadata)
6128                .context("native gateway admission");
6129            let service = gateway_service_error(&error).unwrap();
6130            assert_eq!(service.code, *code);
6131            assert_eq!(service.scope.as_deref(), Some("app"));
6132            assert_eq!(service.request_id.as_deref(), Some("request-1"));
6133            let serialized = serde_json::to_value(service).unwrap();
6134            assert!(serialized.get("providerCost").is_none());
6135            assert_eq!(serialized["documentationUrl"], crate::service_errors::ERROR_DOCUMENTATION_URL);
6136            assert_eq!(gateway_reconnect_delay(1, &error, &mut true), Some(Duration::from_secs(120)));
6137        }
6138    }
6139
6140    #[test]
6141    fn legacy_edge_frames_preserve_canonical_scope_and_owner_cooldown() {
6142        for (legacy, scope) in [("edge-rate-exceeded", "ip"), ("route-rate-exceeded", "avenue")] {
6143            let frame = decode_frame(Message::Text(serde_json::json!({
6144                "type": "error", "code": legacy, "message": "rate limited", "retryable": true,
6145                "retryAfterMs": 120000, "requestId": "legacy-1"
6146            }).to_string().into())).unwrap().unwrap();
6147            let ServerFrame::Error { code, message, retryable, metadata, .. } = frame else {
6148                panic!("expected legacy denial");
6149            };
6150            let error = gateway_server_error_metadata(code, message, retryable, metadata);
6151            let service = gateway_service_error(&error).unwrap();
6152            assert_eq!(service.code, "edge-rate-limited");
6153            assert_eq!(service.scope.as_deref(), Some(scope));
6154            assert_eq!(service.request_id.as_deref(), Some("legacy-1"));
6155            assert_eq!(gateway_reconnect_delay(1, &error, &mut true), Some(Duration::from_secs(120)));
6156        }
6157    }
6158
6159    #[tokio::test]
6160    async fn service_error_observation_is_capability_local_and_preserves_retryability() {
6161        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6162            endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
6163            device_id: "device-1".into(), platform_type: "desktop".into(),
6164            avenue: NativeCoordinationAvenue { kind: "space".into(), id: "space-1".into() },
6165            architecture: None, room_delivery: None,
6166            grant_provider: Arc::new(RecordingGrantProvider::default()),
6167            #[cfg(feature = "managed-group-encryption")]
6168            managed_group_signer: None,
6169            #[cfg(feature = "managed-group-encryption")]
6170            managed_group_store: None,
6171        }).unwrap();
6172        let mut observations = adapter.subscribe_service_errors();
6173        let isolated = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6174            endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
6175            device_id: "device-2".into(), platform_type: "desktop".into(),
6176            avenue: NativeCoordinationAvenue { kind: "space".into(), id: "space-2".into() },
6177            architecture: None, room_delivery: None,
6178            grant_provider: Arc::new(RecordingGrantProvider::default()),
6179            #[cfg(feature = "managed-group-encryption")]
6180            managed_group_signer: None,
6181            #[cfg(feature = "managed-group-encryption")]
6182            managed_group_store: None,
6183        }).unwrap();
6184        let mut isolated_observations = isolated.subscribe_service_errors();
6185        let known = anyhow::Error::new(ControlPlaneHttpError {
6186            status: 429, message: "limited".into(), reason: None,
6187            service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
6188                "code": "app-rate-limited", "retryable": true,
6189                "scope": "app", "operation": "gateway.grant.issue",
6190                "requestId": "request-1", "providerCost": 99,
6191            })),
6192        }).context("wrapped gateway denial");
6193        assert!(adapter.handle_gateway_failure(&known).await);
6194        let observation = observations.recv().await.unwrap();
6195        assert_eq!(observation.avenue.kind, "space");
6196        assert_eq!(observation.avenue.id, "space-1");
6197        assert_eq!(observation.runtime_instance_id, adapter.runtime_instance_id);
6198        assert_eq!(observation.service_error.code, "app-rate-limited");
6199        assert_eq!(observation.service_error.request_id.as_deref(), Some("request-1"));
6200        let serialized = serde_json::to_value(&observation).unwrap();
6201        assert!(serialized.get("avenue").is_some());
6202        assert!(serialized.get("runtimeInstanceId").is_some());
6203        assert!(serialized.get("serviceError").is_some());
6204        assert!(serialized.get("runtime_instance_id").is_none());
6205        assert!(serialized["serviceError"].get("providerCost").is_none());
6206        assert!(serialized["serviceError"].get("message").is_none());
6207        assert!(matches!(isolated_observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
6208        let terminal = anyhow::Error::new(ControlPlaneHttpError {
6209            status: 403, message: "revoked".into(), reason: None,
6210            service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
6211                "code": "app-budget-exhausted", "retryable": false,
6212                "scope": "avenue", "operation": "gateway.grant.issue",
6213            })),
6214        }).context("wrapped terminal denial");
6215        assert!(!adapter.handle_gateway_failure(&terminal).await);
6216        let terminal_observation = observations.recv().await.unwrap();
6217        assert!(!terminal_observation.service_error.retryable);
6218        assert!(matches!(isolated_observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
6219        let unknown = anyhow!("connection reset");
6220        assert!(adapter.handle_gateway_failure(&unknown).await);
6221        assert!(matches!(observations.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
6222        adapter.stop().await;
6223        isolated.stop().await;
6224    }
6225
6226    #[tokio::test]
6227    async fn closed_capability_rejects_service_error_subscription() {
6228        let handle = NativeCapabilityHandle::new(GatewayOptions {
6229            endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
6230            device_id: "device-1".into(), platform_type: "desktop".into(),
6231            avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
6232            architecture: None, room_delivery: None,
6233            grant_provider: Arc::new(RecordingGrantProvider::default()),
6234            #[cfg(feature = "managed-group-encryption")]
6235            managed_group_signer: None,
6236            #[cfg(feature = "managed-group-encryption")]
6237            managed_group_store: None,
6238        }).unwrap();
6239        handle.close().await;
6240        assert!(handle.subscribe_service_errors().is_err());
6241    }
6242
6243    #[tokio::test]
6244    async fn native_actor_republication_cannot_bypass_service_cooldown() {
6245        struct LimitedProvider(std::sync::atomic::AtomicUsize);
6246        #[async_trait]
6247        impl NativeGatewayGrantProvider for LimitedProvider {
6248            async fn grant(&self, _: NativeGatewayGrantRequest) -> Result<NativeGatewayGrant> {
6249                self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
6250                Err(anyhow::Error::new(ControlPlaneHttpError {
6251                    status: 429, message: "try later".into(), reason: None,
6252                    service_error: crate::service_errors::ServiceError::from_value(&serde_json::json!({
6253                        "code": "app-rate-limited", "retryable": true, "retryAfterMs": 120000
6254                    })),
6255                }))
6256            }
6257        }
6258        let provider = Arc::new(LimitedProvider(0.into()));
6259        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6260            endpoint: "https://gateway.example.test".into(), app_tag: "app_native_test".into(),
6261            device_id: "device-1".into(), platform_type: "desktop".into(),
6262            avenue: NativeCoordinationAvenue { kind: "user".into(), id: "principal-1".into() },
6263            architecture: None, room_delivery: None, grant_provider: provider.clone(),
6264            #[cfg(feature = "managed-group-encryption")]
6265            managed_group_signer: None,
6266            #[cfg(feature = "managed-group-encryption")]
6267            managed_group_store: None,
6268        }).unwrap();
6269        tokio::time::pause();
6270        let desired = DesiredPresence {
6271            user_id: "principal-1".into(), local_node_id: "node-1".into(),
6272            ticket: "ticket-1".into(), device_name: "Native".into(),
6273            metadata: None, ttl_ms: 900000, online: true,
6274        };
6275        let (reply, first) = oneshot::channel();
6276        adapter.commands.send(Command::Publish { desired: desired.clone(), reply }).await.unwrap();
6277        for _ in 0..20 { tokio::task::yield_now().await; }
6278        tokio::time::advance(Duration::from_millis(1)).await;
6279        for _ in 0..20 { tokio::task::yield_now().await; }
6280        assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
6281        tokio::time::advance(Duration::from_secs(60)).await;
6282        let (reply, second) = oneshot::channel();
6283        adapter.commands.send(Command::Publish { desired, reply }).await.unwrap();
6284        assert!(first.await.unwrap().is_err(), "old publication is superseded");
6285        for _ in 0..20 { tokio::task::yield_now().await; }
6286        assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
6287        tokio::time::advance(Duration::from_secs(59)).await;
6288        for _ in 0..20 { tokio::task::yield_now().await; }
6289        assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 1);
6290        tokio::time::advance(Duration::from_millis(1001)).await;
6291        for _ in 0..20 { tokio::task::yield_now().await; }
6292        assert_eq!(provider.0.load(std::sync::atomic::Ordering::SeqCst), 2);
6293        adapter.stop().await;
6294        assert!(second.await.unwrap().is_err());
6295    }
6296
6297    #[test]
6298    fn service_cooldown_is_a_minimum_not_capped_network_backoff() {
6299        let denial = |metadata: serde_json::Value| {
6300            anyhow::Error::new(ControlPlaneHttpError {
6301                status: 429,
6302                message: "try later".into(),
6303                reason: None,
6304                service_error: crate::service_errors::ServiceError::from_value(&metadata),
6305            })
6306            .context("obtain native coordination grant")
6307        };
6308        let error = denial(serde_json::json!({
6309            "code": "app-rate-limited", "retryable": true, "retryAfterMs": 120000
6310        }));
6311        let mut allowance = true;
6312        assert_eq!(gateway_service_retry_delay(&error), Duration::from_secs(120));
6313        assert_eq!(gateway_reconnect_delay(1, &error, &mut allowance), Some(Duration::from_secs(120)));
6314        assert_eq!(gateway_reconnect_delay(10, &error, &mut allowance), Some(Duration::from_secs(120)));
6315        assert_eq!(gateway_reconnect_delay(MAX_RECONNECT_ATTEMPTS, &error, &mut allowance), None);
6316        assert!(allowance, "financial denial cannot consume the host return allowance");
6317
6318        let reset = now_ms() + 300000;
6319        let error = denial(serde_json::json!({
6320            "code": "principal-rate-limited", "retryable": true,
6321            "retryAfterMs": 120000, "resetAt": reset
6322        }));
6323        let before = now_ms();
6324        let delay = gateway_service_retry_delay(&error);
6325        assert!(delay >= Duration::from_millis(reset.saturating_sub(now_ms())));
6326        assert!(delay <= Duration::from_millis(reset.saturating_sub(before)));
6327
6328        let terminal = denial(serde_json::json!({
6329            "code": "credit-exhausted", "retryable": false, "retryAfterMs": 120000
6330        }));
6331        assert_eq!(gateway_reconnect_delay(1, &terminal, &mut allowance), None);
6332        let legacy = gateway_connect_error("network outage", true);
6333        assert_eq!(gateway_service_retry_delay(&legacy), Duration::ZERO);
6334        assert_eq!(gateway_reconnect_delay(2, &legacy, &mut allowance), Some(reconnect_delay(2)));
6335    }
6336
6337    #[test]
6338    fn terminal_control_plane_rejections_are_not_retried() {
6339        let revoked = anyhow::Error::new(ControlPlaneHttpError {
6340            service_error: None,
6341            status: 403,
6342            message: "device revoked".into(),
6343            reason: Some("device-certificate-revoked".into()),
6344        })
6345        .context("obtain native coordination grant");
6346        let throttled = anyhow::Error::new(ControlPlaneHttpError {
6347            service_error: None,
6348            status: 429,
6349            message: "try later".into(),
6350            reason: None,
6351        })
6352        .context("obtain native coordination grant");
6353        let unavailable = anyhow::Error::new(ControlPlaneHttpError {
6354            service_error: None,
6355            status: 503,
6356            message: "provider unavailable".into(),
6357            reason: None,
6358        })
6359        .context("obtain native coordination grant");
6360
6361        assert!(!is_retryable_gateway_error(&revoked));
6362        assert!(is_retryable_gateway_error(&throttled));
6363        assert!(is_retryable_gateway_error(&unavailable));
6364    }
6365
6366    #[test]
6367    fn lease_fallback_error_owns_terminal_reconnect_classification() {
6368        let permanent_lease_error = gateway_connect_error("missing operation price", false);
6369        let transient_fallback_error = gateway_connect_error("control plane timeout", true);
6370        let terminal =
6371            terminal_refresh_error(permanent_lease_error, Some(transient_fallback_error));
6372
6373        assert!(is_retryable_gateway_error(&terminal));
6374        assert_eq!(terminal.to_string(), "control plane timeout");
6375    }
6376
6377    #[tokio::test]
6378    async fn native_keepalive() {
6379        let (mut sender, mut receiver) = futures::channel::mpsc::channel(1);
6380        send_socket_keepalive(&mut sender).await.unwrap();
6381        assert!(
6382            matches!(receiver.next().await, Some(Message::Text(payload)) if payload == "ping")
6383        );
6384        assert!(decode_frame(Message::Text("pong".into())).unwrap().is_none());
6385    }
6386
6387    #[cfg(feature = "testing-endpoints")]
6388    #[tokio::test]
6389    async fn busy_socket_keepalive() {
6390        // Exercise the actual actor, not a second select loop in a test.
6391        // Incoming frames must not restart its keepalive deadline.
6392        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
6393        let endpoint = format!("http://{}", listener.local_addr().unwrap());
6394        let provider = Arc::new(RecordingGrantProvider {
6395            gateway_url: Some(endpoint.clone()),
6396            token: Some(format!("local.{}.test",
6397                base64::engine::general_purpose::URL_SAFE_NO_PAD
6398                    .encode(br#"{"jti":"busy-socket-grant"}"#))),
6399            ..Default::default()
6400        });
6401        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6402            endpoint,
6403            app_tag: "app_native_test".into(),
6404            device_id: "busy-device".into(),
6405            platform_type: "native".into(),
6406            avenue: NativeCoordinationAvenue { kind: "user".into(), id: "busy-user".into() },
6407            architecture: None,
6408            room_delivery: None,
6409            grant_provider: provider.clone(),
6410            #[cfg(feature = "managed-group-encryption")]
6411            managed_group_signer: None,
6412            #[cfg(feature = "managed-group-encryption")]
6413            managed_group_store: None,
6414        }).unwrap();
6415        let server = tokio::spawn(async move {
6416            let (tcp, _) = listener.accept().await.unwrap();
6417            let mut socket = tokio_tungstenite::accept_hdr_async(tcp,
6418                |_: &tokio_tungstenite::tungstenite::handshake::server::Request,
6419                 mut response: tokio_tungstenite::tungstenite::handshake::server::Response| {
6420                    response.headers_mut().insert("Sec-WebSocket-Protocol",
6421                        HeaderValue::from_static(GATEWAY_PROTOCOL));
6422                    Ok(response)
6423                }).await.unwrap();
6424            let auth: serde_json::Value = serde_json::from_str(
6425                socket.next().await.unwrap().unwrap().to_text().unwrap()).unwrap();
6426            assert_eq!(auth["socketLiveness"], "ping-v1");
6427            socket.send(Message::Text(
6428                serde_json::json!({"type":"ready", "budgetRemainingMicrousd":1000})
6429                    .to_string().into())).await.unwrap();
6430            let start = tokio::time::Instant::now();
6431            let mut traffic = tokio::time::interval(Duration::from_millis(100));
6432            let mut pings = 0;
6433            loop {
6434                tokio::select! {
6435                    _ = traffic.tick() => {
6436                        socket.send(Message::Text("pong".into())).await.unwrap();
6437                    }
6438                    frame = socket.next() => {
6439                        assert!(matches!(frame, Some(Ok(Message::Text(payload))) if payload == "ping"),
6440                            "keepalive must not republish presence or refresh authorization");
6441                        pings += 1;
6442                        if pings >= 2 && start.elapsed() >= SOCKET_KEEPALIVE_INTERVAL / 2 {
6443                            return;
6444                        }
6445                    }
6446                }
6447            }
6448        });
6449        let published = tokio::time::timeout(Duration::from_secs(3), adapter.update_presence(
6450            "busy-user", &iroh::SecretKey::from_bytes(&[19; 32]).public().to_string(),
6451            "busy-ticket", true, "Busy", 900_000, None)).await;
6452        let mut server = server;
6453        let observed = tokio::time::timeout(SOCKET_KEEPALIVE_INTERVAL * 2 + Duration::from_secs(5),
6454            &mut server).await;
6455        adapter.stop().await;
6456        if observed.is_err() {
6457            server.abort();
6458            let _ = server.await;
6459        }
6460        published.expect("initial presence timed out").unwrap();
6461        observed.expect("incoming traffic starved native keepalive").unwrap();
6462        assert_eq!(provider.requests.lock().unwrap().len(), 1,
6463            "keepalive must reuse the authenticated socket and grant");
6464    }
6465
6466    #[test]
6467    fn only_local_device_delete_terminates_the_gateway_presence_owner() {
6468        let (local_reply, _local_result) = oneshot::channel();
6469        let local = Command::Delete {
6470            user_id: "user-1".into(),
6471            device_id: "local-device".into(),
6472            reply: local_reply,
6473        };
6474        let (remote_reply, _remote_result) = oneshot::channel();
6475        let remote = Command::Delete {
6476            user_id: "user-1".into(),
6477            device_id: "remote-device".into(),
6478            reply: remote_reply,
6479        };
6480        assert!(command_deletes_device(&local, "local-device"));
6481        assert!(!command_deletes_device(&remote, "local-device"));
6482    }
6483
6484    #[tokio::test]
6485    async fn healthy_lease_response_advances_the_local_device_projection() {
6486        let provider = Arc::new(RecordingGrantProvider::default());
6487        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6488            endpoint: "https://gateway.example.test".to_string(),
6489            app_tag: "app_native_test".to_string(),
6490            device_id: "device-1".to_string(),
6491            platform_type: "desktop".to_string(),
6492            avenue: NativeCoordinationAvenue {
6493                kind: "user".to_string(),
6494                id: "principal-1".to_string(),
6495            },
6496            architecture: None,
6497            room_delivery: None,
6498            grant_provider: provider,
6499            #[cfg(feature = "managed-group-encryption")]
6500            managed_group_signer: None,
6501            #[cfg(feature = "managed-group-encryption")]
6502            managed_group_store: None,
6503        })
6504        .expect("adapter");
6505        adapter.shared.devices.write().await.insert(
6506            "device-1".to_string(),
6507            Device {
6508                app_tag: Some("app_native_test".to_string()),
6509                device_id: "device-1".to_string(),
6510                user_id: Some("principal-1".to_string()),
6511                device_name: "Native".to_string(),
6512                platform_type: Some("desktop".to_string()),
6513                capabilities: None,
6514                session_id: Some("runtime-1".to_string()),
6515                node_id: Some("00".repeat(32)),
6516                tag: None,
6517                kind: None,
6518                metadata: None,
6519                online: true,
6520                ticket: Some("ticket-1".to_string()),
6521                last_seen_at: None,
6522                expires_at: Some(serde_json::json!(1_000_u64)),
6523                created_at: None,
6524                updated_at: None,
6525                excluded_peers: Vec::new(),
6526            },
6527        );
6528
6529        update_local_device_expiry(&adapter, 2_000).await;
6530
6531        let devices = adapter.shared.devices.read().await;
6532        assert_eq!(
6533            devices
6534                .get("device-1")
6535                .and_then(|device| device.expires_at.clone()),
6536            Some(serde_json::json!(2_000_u64)),
6537        );
6538    }
6539
6540    #[tokio::test]
6541    async fn native_membership_pages_publish_only_after_the_complete_snapshot() {
6542        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6543            endpoint: "https://gateway.example.test".to_string(),
6544            app_tag: "app_native_test".to_string(),
6545            device_id: "device-local".to_string(),
6546            platform_type: "desktop".to_string(),
6547            avenue: NativeCoordinationAvenue {
6548                kind: "room".to_string(),
6549                id: "room-paged".to_string(),
6550            },
6551            architecture: Some(RoomArchitectureMode::Managed),
6552            room_delivery: Some(RoomDelivery::Reliable),
6553            grant_provider: Arc::new(RecordingGrantProvider::default()),
6554            #[cfg(feature = "managed-group-encryption")]
6555            managed_group_signer: None,
6556            #[cfg(feature = "managed-group-encryption")]
6557            managed_group_store: None,
6558        })
6559        .expect("adapter");
6560        let members = (0..201)
6561            .map(|index| GatewayMember {
6562                user_id: Some(format!("user-{index}")),
6563                device_id: format!("device-{index:03}"),
6564                device_name: format!("Device {index}"),
6565                platform_type: "test".to_string(),
6566                metadata: None,
6567                capabilities: None,
6568                online: true,
6569                updated_at_ms: 1_000,
6570                expires_at_ms: 60_000,
6571            })
6572            .collect::<Vec<_>>();
6573
6574        for page_index in 0..3 {
6575            handle_frame(
6576                &adapter,
6577                ServerFrame::MembershipPage {
6578                    snapshot_id: "connection:paged-room".to_string(),
6579                    page_index,
6580                    page_count: 3,
6581                    member_count: members.len(),
6582                    members: members[page_index * 100..((page_index + 1) * 100).min(members.len())]
6583                        .to_vec(),
6584                },
6585            )
6586            .await
6587            .expect("membership page");
6588            if page_index < 2 {
6589                assert!(adapter.shared.members.read().await.is_empty());
6590            }
6591        }
6592
6593        assert_eq!(adapter.shared.members.read().await.len(), 201);
6594        assert!(adapter
6595            .shared
6596            .membership_pages
6597            .lock()
6598            .expect("membership page lock")
6599            .is_none());
6600    }
6601
6602    #[tokio::test]
6603    async fn v2_topology_projects_room_architecture_and_selected_routes() {
6604        let provider = Arc::new(RecordingGrantProvider::default());
6605        let adapter = NativeCoordinationGatewaySignaling::new(GatewayOptions {
6606            endpoint: "https://gateway.example.test".to_string(),
6607            app_tag: "app_native_test".to_string(),
6608            device_id: "device-local".to_string(),
6609            platform_type: "desktop".to_string(),
6610            avenue: NativeCoordinationAvenue {
6611                kind: "room".to_string(),
6612                id: "room-adaptive".to_string(),
6613            },
6614            architecture: Some(RoomArchitectureMode::Auto),
6615            room_delivery: Some(RoomDelivery::LatestState),
6616            grant_provider: provider,
6617            #[cfg(feature = "managed-group-encryption")]
6618            managed_group_signer: None,
6619            #[cfg(feature = "managed-group-encryption")]
6620            managed_group_store: None,
6621        })
6622        .expect("adapter");
6623        let now = now_ms();
6624        adapter.shared.members.write().await.insert(
6625            "device-peer".to_string(),
6626            GatewayMember {
6627                user_id: None,
6628                device_id: "device-peer".to_string(),
6629                device_name: "Peer".to_string(),
6630                platform_type: "desktop".to_string(),
6631                metadata: None,
6632                capabilities: None,
6633                online: true,
6634                updated_at_ms: now as i64,
6635                expires_at_ms: now.saturating_add(60_000) as i64,
6636            },
6637        );
6638        *adapter.shared.active_grant.write().await =
6639            Some(("grant-current".to_string(), now.saturating_add(60_000)));
6640
6641        let current_lease = GatewayTopologyLease {
6642            schema_version: 2,
6643            admission_peers: None,
6644            topology_revision: 1,
6645            previous_topology_revision: None,
6646            avenue: adapter.avenue.clone(),
6647            grant_jti: "grant-current".to_string(),
6648            expires_at_ms: now.saturating_add(30_000),
6649            architecture: GatewayArchitectureLease {
6650                policy_version: "room-architecture-v1".to_string(),
6651                epoch: 1,
6652                previous_architecture_epoch: None,
6653                requested_mode: RoomArchitectureMode::Auto,
6654                effective_mode: EffectiveRoomArchitecture::Sparse,
6655                phase: RoomArchitecturePhase::Settled,
6656                reason: RoomArchitectureReason::Size,
6657                held_credits_microusd: 25_000,
6658                quote_expires_at_ms: now.saturating_add(30_000),
6659                reservation_id: Some("reservation:test".to_string()),
6660            },
6661            group_encryption: None,
6662            active: vec![GatewayDevice {
6663                user_id: None,
6664                device_id: "device-peer".to_string(),
6665                // This fixture tests a sparse lease, not authority admission.
6666                runtime_instance_id: "runtime-peer".to_string(),
6667                node_id: "node-peer".to_string(),
6668                device_name: "Peer".to_string(),
6669                platform_type: "desktop".to_string(),
6670                ticket: "ticket-peer".to_string(),
6671                metadata: None,
6672                capabilities: None,
6673                excluded_peers: Vec::new(),
6674                online: true,
6675                updated_at_ms: now as i64,
6676                expires_at_ms: now.saturating_add(60_000) as i64,
6677            }],
6678            backups: Vec::new(),
6679        };
6680        let mut wrong_grant = current_lease.clone();
6681        wrong_grant.grant_jti = "grant-forged".to_string();
6682        assert!(accept_topology_lease(&adapter, wrong_grant).await.is_err());
6683        assert_eq!(adapter.room_architecture(), None);
6684
6685        accept_topology_lease(&adapter, current_lease.clone())
6686            .await
6687            .expect("v2 topology lease");
6688
6689        assert_eq!(
6690            adapter.room_architecture(),
6691            Some(RoomArchitectureSnapshot {
6692                requested: RoomArchitectureMode::Auto,
6693                effective: EffectiveRoomArchitecture::Sparse,
6694                epoch: 1,
6695                phase: RoomArchitecturePhase::Settled,
6696                reason: RoomArchitectureReason::Size,
6697                held_credits_usd: 0.025,
6698                quote_expires_at_ms: now.saturating_add(30_000),
6699            }),
6700        );
6701        assert_eq!(
6702            adapter
6703                .shared
6704                .devices
6705                .read()
6706                .await
6707                .get("device-peer")
6708                .and_then(|device| device.node_id.as_deref()),
6709            Some("node-peer"),
6710        );
6711
6712        accept_topology_lease(&adapter, current_lease.clone())
6713            .await
6714            .expect("exact replay is stale and idempotent");
6715        assert_eq!(
6716            adapter.room_architecture().map(|snapshot| snapshot.epoch),
6717            Some(1),
6718        );
6719        let mut forged_replay = current_lease;
6720        forged_replay.grant_jti = "grant-forged".to_string();
6721        assert!(accept_topology_lease(&adapter, forged_replay)
6722            .await
6723            .is_err());
6724        assert_eq!(
6725            adapter.room_architecture().map(|snapshot| snapshot.epoch),
6726            Some(1),
6727        );
6728
6729        let wrong_predecessor = GatewayTopologyLease {
6730            schema_version: 2,
6731            admission_peers: None,
6732            topology_revision: 2,
6733            previous_topology_revision: Some(99),
6734            avenue: adapter.avenue.clone(),
6735            grant_jti: "grant-current".to_string(),
6736            expires_at_ms: now.saturating_add(30_000),
6737            architecture: GatewayArchitectureLease {
6738                policy_version: "room-architecture-v1".to_string(),
6739                epoch: 2,
6740                previous_architecture_epoch: Some(1),
6741                requested_mode: RoomArchitectureMode::Auto,
6742                effective_mode: EffectiveRoomArchitecture::Managed,
6743                phase: RoomArchitecturePhase::Preparing,
6744                reason: RoomArchitectureReason::Size,
6745                held_credits_microusd: 25_000,
6746                quote_expires_at_ms: now.saturating_add(30_000),
6747                reservation_id: Some("reservation:managed".to_string()),
6748            },
6749            group_encryption: None,
6750            active: Vec::new(),
6751            backups: Vec::new(),
6752        };
6753        assert!(accept_topology_lease(&adapter, wrong_predecessor)
6754            .await
6755            .is_err());
6756        assert_eq!(
6757            adapter.room_architecture().map(|snapshot| snapshot.epoch),
6758            Some(1),
6759        );
6760
6761        accept_topology_lease(
6762            &adapter,
6763            GatewayTopologyLease {
6764                schema_version: 2,
6765                admission_peers: None,
6766                topology_revision: 2,
6767                previous_topology_revision: Some(1),
6768                avenue: adapter.avenue.clone(),
6769                grant_jti: "grant-current".to_string(),
6770                expires_at_ms: now.saturating_add(30_000),
6771                architecture: GatewayArchitectureLease {
6772                    policy_version: "room-architecture-v1".to_string(),
6773                    epoch: 2,
6774                    previous_architecture_epoch: Some(1),
6775                    requested_mode: RoomArchitectureMode::Auto,
6776                    effective_mode: EffectiveRoomArchitecture::Managed,
6777                    phase: RoomArchitecturePhase::Preparing,
6778                    reason: RoomArchitectureReason::Size,
6779                    held_credits_microusd: 25_000,
6780                    quote_expires_at_ms: now.saturating_add(30_000),
6781                    reservation_id: Some("reservation:managed".to_string()),
6782                },
6783                group_encryption: None,
6784                active: vec![GatewayDevice {
6785                    user_id: None,
6786                    device_id: "device-peer".to_string(),
6787                    runtime_instance_id: "runtime-peer".to_string(),
6788                    node_id: "node-peer".to_string(),
6789                    device_name: "Peer".to_string(),
6790                    platform_type: "desktop".to_string(),
6791                    ticket: "ticket-peer".to_string(),
6792                    metadata: None,
6793                    capabilities: None,
6794                    excluded_peers: Vec::new(),
6795                    online: true,
6796                    updated_at_ms: now as i64,
6797                    expires_at_ms: now.saturating_add(60_000) as i64,
6798                }],
6799                backups: Vec::new(),
6800            },
6801        )
6802        .await
6803        .expect("preparing managed lease keeps sparse routes without a group lease");
6804        assert_eq!(
6805            adapter
6806                .room_architecture()
6807                .map(|snapshot| (snapshot.effective, snapshot.phase)),
6808            Some((
6809                EffectiveRoomArchitecture::Managed,
6810                RoomArchitecturePhase::Preparing,
6811            )),
6812        );
6813        assert!(adapter
6814            .shared
6815            .devices
6816            .read()
6817            .await
6818            .contains_key("device-peer"));
6819    }
6820
6821    #[tokio::test]
6822    async fn command_completion_preserves_permanent_error_classification() {
6823        let (reply, result) = oneshot::channel();
6824        let actor_error = complete_command::<()>(
6825            reply,
6826            Err(gateway_connect_error("invalid payload", false)),
6827            "presence publication",
6828        )
6829        .expect_err("permanent operation must fail");
6830
6831        assert!(!is_retryable_gateway_error(&actor_error));
6832        assert_eq!(
6833            result
6834                .await
6835                .expect("caller receives reply")
6836                .expect_err("caller receives operation failure")
6837                .to_string(),
6838            "invalid payload"
6839        );
6840    }
6841
6842    #[test]
6843    fn optional_device_fields_are_omitted_instead_of_null() {
6844        let value = serde_json::to_value(OutboundGatewayDevice {
6845            device_id: "device:test".into(),
6846            runtime_instance_id: "runtime:test".into(),
6847            node_id: "00".repeat(32),
6848            device_name: "Test".into(),
6849            platform_type: "macos".into(),
6850            ticket: "ticket".into(),
6851            metadata: None,
6852            capabilities: None,
6853            excluded_peers: Vec::new(),
6854            online: true,
6855        })
6856        .unwrap();
6857        assert!(value.get("metadata").is_none());
6858        assert!(value.get("capabilities").is_none());
6859    }
6860}