1#[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;
53const 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
90fn 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#[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
235pub 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#[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#[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
362pub 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 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 #[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
475pub 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 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 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 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
1122struct AuthorityPeerBinding {
1126 client: std::sync::Weak<crate::Client>,
1127 scope: String,
1128 revision: u64,
1129 admission_revision: u64,
1131 last_payload: String,
1132 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 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 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 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 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 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 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 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 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 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 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 *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 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
2322async 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 *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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 *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 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 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}