1use std::{
2 collections::{BTreeMap, HashMap, HashSet},
3 error::Error,
4 fmt,
5 sync::{Arc, Mutex, MutexGuard, RwLock, RwLockReadGuard, RwLockWriteGuard},
6 time::Duration,
7};
8
9use subc_control::{ClientControlResponse, RouteCloseReason};
10use subc_protocol::{
11 manifest::Concurrency,
12 session::{LiveRoot, ModuleControlResponse, ModuleControlResponseToModule},
13 ErrorBody, Flags, FrameType, Principal, Priority,
14};
15use tokio::sync::{oneshot, Semaphore};
16use tokio::time::Instant;
17use tracing::{debug, info, warn};
18
19use crate::{
20 control::{RouteBindBreakers, RouteBindConcurrency},
21 observability::DaemonCounters,
22 registry::ConnectionId,
23 router::FrameSink,
24 scopes::{BoundScope, ScopeDrain, ScopeTag, ScopeTagChange},
25 Frame, ProjectRootId,
26};
27
28const DEFAULT_MODULE_MANAGED_WINDOW: usize = 32;
30
31const STATELESS_PARALLEL_WINDOW: usize = 1024;
33
34const HEALTH_PROBE_TOMBSTONE_TTL: Duration = Duration::from_secs(5 * 60);
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
44pub struct ModuleEndpointId {
45 pub connection_id: ConnectionId,
46 pub generation: u64,
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
51pub(crate) struct ClientRouteKey {
52 pub connection_id: ConnectionId,
53 pub channel: u16,
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
58pub(crate) struct ModuleRouteKey {
59 pub endpoint: ModuleEndpointId,
60 pub channel: u16,
61}
62
63#[derive(Debug)]
64pub(crate) struct RouteBinding {
65 pub client_connection_id: ConnectionId,
66 pub client_sink: FrameSink,
67 pub client_negotiated_ver: u8,
68 pub client_channel: u16,
69 pub client_epoch: u32,
70 pub module_id: String,
71 pub module_endpoint: ModuleEndpointId,
72 pub module_sink: FrameSink,
73 pub module_negotiated_ver: u8,
74 pub module_channel: u16,
75 pub module_epoch: u32,
76 pub principal: Principal,
77 pub project_root: Option<ProjectRootId>,
78 pub bound_at: Instant,
79 pub flow: Arc<ChannelFlow>,
80 pub scope: Option<BoundScope>,
83}
84
85#[derive(Debug, Clone)]
86pub(crate) enum DataRoute {
87 Client(DataRouteState),
88 Module(DataRouteState),
89}
90
91#[derive(Debug, Clone)]
92pub(crate) enum DataRouteState {
93 Bound(Arc<RouteBinding>),
94 Reserved,
95 EpochMismatch,
96 Absent,
97}
98
99#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128pub(crate) enum GoodbyeTargetKind {
129 Client,
130 Module,
131}
132
133#[derive(Debug, Clone)]
134pub(crate) struct GoodbyeTarget {
135 pub connection_id: ConnectionId,
136 pub sink: FrameSink,
137 pub negotiated_ver: u8,
138 pub channel: u16,
139 pub epoch: u32,
140 pub kind: GoodbyeTargetKind,
141 pub module_id: Option<String>,
146}
147
148#[derive(Debug, Clone, Copy)]
151pub(crate) struct UndeliveredFrame<'a> {
152 pub module_id: Option<&'a str>,
154 pub sink: &'a FrameSink,
156}
157
158fn principal_label(principal: &Principal) -> String {
160 match principal {
161 Principal::Reserved { module_id } => format!("reserved:{module_id}"),
162 Principal::Direct => "direct".to_string(),
163 other => format!("{other:?}"),
164 }
165}
166
167fn connection_principals_locked(inner: &ForwardingInner, connection_id: ConnectionId) -> String {
170 let labels = inner
171 .client_to_module
172 .iter()
173 .filter(|(key, _)| key.connection_id == connection_id)
174 .map(|(_, route)| principal_label(&route.principal))
175 .collect::<std::collections::BTreeSet<_>>();
176 if labels.is_empty() {
177 "none".to_string()
178 } else {
179 labels.into_iter().collect::<Vec<_>>().join(",")
180 }
181}
182
183impl GoodbyeTarget {
184 pub(crate) fn close_on_delivery_failure(&self) -> bool {
187 matches!(self.kind, GoodbyeTargetKind::Client)
188 }
189}
190
191pub(crate) const LATE_MODULE_GOODBYE_DEADLINE: Duration = crate::supervise::DEFAULT_DRAIN_TIMEOUT;
204
205pub(crate) fn send_module_route_goodbye(
217 counters: &DaemonCounters,
218 sink: &FrameSink,
219 frame: Frame,
220 module_id: Option<&str>,
221 context: &'static str,
222) {
223 let channel = frame.header.channel;
224 let epoch = frame.header.epoch;
225 let Err(err) = sink.try_send(frame.clone()) else {
226 return;
227 };
228 let runtime = match tokio::runtime::Handle::try_current() {
231 Ok(runtime) if !sink.is_closed() => runtime,
232 _ => {
233 counters.increment_goodbye_relay_module_dropped(module_id);
234 warn!(
235 module_id = module_id.unwrap_or("unknown"),
236 route_channel = channel,
237 route_epoch = epoch,
238 error = %err,
239 context,
240 "route GOODBYE to module dropped: module connection is closed; not closing shared module connection"
241 );
242 return;
243 }
244 };
245 debug!(
246 module_id = module_id.unwrap_or("unknown"),
247 route_channel = channel,
248 route_epoch = epoch,
249 error = %err,
250 context,
251 "module egress queue refused route GOODBYE; delivering it once the module frees room"
252 );
253 let counters = counters.clone();
254 let sink = sink.clone();
255 let module_id = module_id.map(str::to_string);
256 runtime.spawn(async move {
257 let outcome = tokio::time::timeout(LATE_MODULE_GOODBYE_DEADLINE, sink.send(frame)).await;
258 let why = match outcome {
259 Ok(Ok(())) => {
260 debug!(
261 module_id = module_id.as_deref().unwrap_or("unknown"),
262 route_channel = channel,
263 route_epoch = epoch,
264 context,
265 "late route GOODBYE delivered to module"
266 );
267 return;
268 }
269 Ok(Err(err)) => err.to_string(),
270 Err(_) => format!(
271 "module egress queue had no room within {LATE_MODULE_GOODBYE_DEADLINE:?}"
272 ),
273 };
274 counters.increment_goodbye_relay_module_dropped(module_id.as_deref());
275 warn!(
276 module_id = module_id.as_deref().unwrap_or("unknown"),
277 route_channel = channel,
278 route_epoch = epoch,
279 error = %why,
280 context,
281 "route GOODBYE to module dropped under backpressure; not closing shared module connection"
282 );
283 });
284}
285
286#[derive(Debug, Clone)]
292pub(crate) struct EndpointRoute {
293 pub goodbye_target: GoodbyeTarget,
294 pub principal: Principal,
295 pub bound_at: Instant,
296 pub draining: bool,
297 pub drain_reason: Option<RouteCloseReason>,
301}
302
303#[derive(Debug, Clone)]
306pub(crate) struct ScopeDrainedRoute {
307 pub reason: RouteCloseReason,
308 pub module_id: String,
309 pub client: GoodbyeTarget,
310 pub module: GoodbyeTarget,
311}
312
313#[derive(Debug)]
314pub(crate) struct PendingRouteBindRelay {
315 pub endpoint: ModuleEndpointId,
316 pub module_sink: FrameSink,
317 pub negotiated_ver: u8,
318 pub client_channel: u16,
319 pub client_epoch: u32,
320 pub module_channel: u16,
321 pub module_epoch: u32,
322 pub corr: u64,
323 pub receiver: oneshot::Receiver<RouteBindRelayOutcome>,
324}
325
326#[derive(Debug, Clone)]
327pub(crate) struct ModuleDrainTarget {
328 pub endpoint: ModuleEndpointId,
329 pub sink: FrameSink,
330 pub negotiated_ver: u8,
331 pub abandoned_bindings: Vec<GoodbyeTarget>,
332 pub excluded_subscriptions: u32,
333}
334
335#[cfg(unix)]
338#[derive(Debug, Clone)]
339pub(crate) struct ModuleConnectionTarget {
340 pub module_id: String,
341 pub endpoint: ModuleEndpointId,
342 pub sink: FrameSink,
343 pub negotiated_ver: u8,
344}
345
346#[derive(Debug, Clone)]
347pub(crate) enum RouteBindRelayOutcome {
348 Accepted,
349 Rejected(ErrorBody),
350 ModuleGone(String),
351}
352
353#[derive(Debug, Clone, Copy, PartialEq, Eq)]
355pub(crate) struct ForwardingCutover {
356 pub promoted: ModuleEndpointId,
358 pub incumbent: Option<ModuleEndpointId>,
361}
362
363#[derive(Debug)]
365pub(crate) struct ConnectionCleanup {
366 pub released: Vec<GoodbyeTarget>,
368 pub abandoned_relays: u32,
371}
372
373#[derive(Debug, Clone)]
374pub(crate) struct PendingRelayCompletion {
375 pub settled: bool,
376 pub abandoned: Option<GoodbyeTarget>,
377}
378
379#[derive(Debug)]
380pub(crate) struct PendingModuleControlRpc {
381 pub endpoint: ModuleEndpointId,
382 pub module_sink: FrameSink,
383 pub negotiated_ver: u8,
384 pub corr: u64,
385 pub receiver: oneshot::Receiver<ModuleControlRpcOutcome>,
386}
387
388#[derive(Debug, Clone)]
389pub(crate) enum ModuleControlRpcOutcome {
390 Response(ModuleControlResponse),
391 Rejected(ErrorBody),
392 ModuleGone(String),
393 MalformedResponse(String),
394 UnexpectedOp { expected: String, actual: String },
395 DeadlineElapsed,
396}
397
398#[derive(Debug, Clone, PartialEq, Eq)]
399pub(crate) enum ModuleControlRpcCompletion {
400 Unknown,
401 Settled,
402 LateHealthAnswer {
403 module_id: String,
404 latency: Duration,
405 },
406}
407
408#[derive(Debug)]
409struct PendingModuleControlRpcEntry {
410 expected_op: String,
411 deadline: Instant,
412 health_probe_started_at: Option<Instant>,
413 sender: oneshot::Sender<ModuleControlRpcOutcome>,
414}
415
416#[derive(Debug)]
417struct HealthProbeTombstone {
418 expected_op: String,
419 module_id: String,
420 probe_started_at: Instant,
421 expires_at: Instant,
422}
423
424#[derive(Debug, Clone)]
425struct RouteReservation {
426 client_key: ClientRouteKey,
427 module_key: ModuleRouteKey,
428 client_epoch: u32,
429 module_epoch: u32,
430 project_root: Option<ProjectRootId>,
431}
432
433#[derive(Debug)]
434struct PendingRouteBindRelayEntry {
435 reservation: RouteReservation,
436 client_sink: FrameSink,
437 client_negotiated_ver: u8,
438 client_permit: crate::router::EgressPermit,
439 route_open_frame: Frame,
440 principal: Principal,
441 scope: Option<BoundScope>,
444 deadline: Instant,
445 relay_enqueued: bool,
446 sender: oneshot::Sender<RouteBindRelayOutcome>,
447}
448
449#[derive(Debug, Clone)]
450pub(crate) enum RouteRelease {
451 Removed(GoodbyeTarget),
452 Stale,
453 Absent,
454}
455
456#[derive(Debug, Clone)]
457pub(crate) enum RoutePollSnapshot {
458 Bound {
459 module_id: String,
460 status: Option<String>,
461 },
462 Absent,
463}
464
465#[derive(Debug, Clone)]
466struct ModuleConnection {
467 endpoint: ModuleEndpointId,
468 sink: FrameSink,
469 negotiated_ver: u8,
470 concurrency: Concurrency,
471}
472
473#[derive(Debug, Default)]
474struct ForwardingInner {
475 operator_confirms: Arc<crate::operator_confirm::OperatorConfirms>,
476 daemon_draining: bool,
477 modules_by_id: HashMap<String, ModuleConnection>,
481 candidates_by_id: HashMap<String, ModuleConnection>,
487 superseded_endpoints: HashMap<ModuleEndpointId, ModuleConnection>,
493 endpoint_by_connection: HashMap<ConnectionId, ModuleEndpointId>,
494 module_id_by_endpoint: HashMap<ModuleEndpointId, String>,
495 draining_endpoints: HashMap<ModuleEndpointId, RouteCloseReason>,
499 closing_connections: HashSet<ConnectionId>,
500 next_generation: u64,
501 reserved_client: HashMap<ClientRouteKey, ModuleRouteKey>,
502 reserved_module: HashMap<ModuleRouteKey, ClientRouteKey>,
503 next_client_channel: HashMap<ConnectionId, u16>,
504 next_module_channel: HashMap<ModuleEndpointId, u16>,
505 client_slot_epochs: HashMap<ClientRouteKey, u32>,
506 module_slot_epochs: HashMap<ModuleRouteKey, u32>,
507 last_published_epoch: HashMap<ClientRouteKey, u32>,
508 client_to_module: HashMap<ClientRouteKey, Arc<RouteBinding>>,
509 module_to_client: HashMap<ModuleRouteKey, Arc<RouteBinding>>,
510 status: HashMap<(ClientRouteKey, u32), String>,
511 pending_relays: HashMap<(ModuleEndpointId, u64), PendingRouteBindRelayEntry>,
512 scope_tags: HashMap<(String, String), ScopeTag>,
517 next_control_corr: HashMap<ModuleEndpointId, u64>,
518 pending_control_rpcs: HashMap<(ModuleEndpointId, u64), PendingModuleControlRpcEntry>,
519 health_probe_tombstones: HashMap<(ModuleEndpointId, u64), HealthProbeTombstone>,
520}
521
522#[derive(Debug, Clone)]
523pub(crate) struct CloseReason {
524 code: &'static str,
525 message: String,
526}
527
528impl CloseReason {
529 pub(crate) fn new(code: &'static str, message: impl Into<String>) -> Self {
530 Self {
531 code,
532 message: message.into(),
533 }
534 }
535}
536
537impl fmt::Display for CloseReason {
538 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
539 write!(f, "{}: {}", self.code, self.message)
540 }
541}
542
543pub(crate) type ConnectionCloseReceiver = oneshot::Receiver<CloseReason>;
544
545#[derive(Debug, Default)]
547pub struct ForwardingTable {
548 inner: Arc<RwLock<ForwardingInner>>,
549 close_registry: Mutex<HashMap<ConnectionId, oneshot::Sender<CloseReason>>>,
550 counters: DaemonCounters,
551 route_bind_breakers: RouteBindBreakers,
556 route_bind_concurrency: RouteBindConcurrency,
559 route_outages: Arc<crate::route_outage::RouteOutageTracker>,
563}
564
565impl ForwardingTable {
566 pub(crate) fn operator_confirms(&self) -> Arc<crate::operator_confirm::OperatorConfirms> {
567 Arc::clone(
568 &self
569 .inner
570 .read()
571 .unwrap_or_else(|p| p.into_inner())
572 .operator_confirms,
573 )
574 }
575
576 #[cfg(test)]
577 pub(crate) fn inject_operator_principal(&self, key: ModuleRouteKey, principal: Principal) {
578 let mut inner = self.write_inner().unwrap();
579 let old = inner.module_to_client.remove(&key).unwrap();
580 let client = ClientRouteKey {
581 connection_id: old.client_connection_id,
582 channel: old.client_channel,
583 };
584 inner.client_to_module.remove(&client);
585 let mut binding = Arc::try_unwrap(old).expect("test binding has no outstanding readers");
586 binding.principal = principal;
587 let binding = Arc::new(binding);
588 inner.client_to_module.insert(client, Arc::clone(&binding));
589 inner.module_to_client.insert(key, binding);
590 }
591
592 pub(crate) fn with_operator_route<T>(
596 &self,
597 connection_id: ConnectionId,
598 channel: u16,
599 epoch: u32,
600 admit: impl FnOnce(Option<&RouteBinding>) -> T,
601 ) -> Result<T, ForwardingError> {
602 let inner = self.read_inner()?;
603 let binding = inner
604 .endpoint_by_connection
605 .get(&connection_id)
606 .and_then(|endpoint| {
607 inner.module_to_client.get(&ModuleRouteKey {
608 endpoint: *endpoint,
609 channel,
610 })
611 })
612 .filter(|binding| binding.module_epoch == epoch);
613 Ok(admit(binding.map(Arc::as_ref)))
614 }
615
616 pub(crate) fn counters(&self) -> DaemonCounters {
617 self.counters.clone()
618 }
619
620 pub(crate) fn route_bind_breakers(&self) -> RouteBindBreakers {
621 self.route_bind_breakers.clone()
622 }
623
624 pub(crate) fn route_bind_concurrency(&self) -> RouteBindConcurrency {
625 self.route_bind_concurrency.clone()
626 }
627
628 pub(crate) fn route_outages(&self) -> Arc<crate::route_outage::RouteOutageTracker> {
629 Arc::clone(&self.route_outages)
630 }
631
632 pub(crate) fn register_connection_close(
633 &self,
634 connection_id: ConnectionId,
635 ) -> ConnectionCloseReceiver {
636 let (sender, receiver) = oneshot::channel();
637 let replaced = self
638 .lock_close_registry()
639 .insert(connection_id, sender)
640 .is_some();
641 if replaced {
642 warn!(
643 connection_id = connection_id.get(),
644 "replaced existing connection close registration"
645 );
646 }
647 receiver
648 }
649
650 pub(crate) fn unregister_connection_close(&self, connection_id: ConnectionId) {
651 self.lock_close_registry().remove(&connection_id);
652 }
653
654 #[cfg(unix)]
663 pub(crate) fn close_all_connections(&self, reason: &CloseReason) -> usize {
664 let senders: Vec<_> = self.lock_close_registry().drain().collect();
665 let count = senders.len();
666 for (_, sender) in senders {
667 let _ = sender.send(reason.clone());
668 }
669 count
670 }
671
672 #[cfg(unix)]
676 pub(crate) fn module_connections(
677 &self,
678 ) -> Result<Vec<ModuleConnectionTarget>, ForwardingError> {
679 let inner = self.read_inner()?;
680 let mut seen = HashSet::new();
681 Ok(inner
682 .modules_by_id
683 .values()
684 .chain(inner.candidates_by_id.values())
685 .chain(inner.superseded_endpoints.values())
686 .filter(|module| seen.insert(module.endpoint))
687 .map(|module| ModuleConnectionTarget {
688 module_id: inner
689 .module_id_by_endpoint
690 .get(&module.endpoint)
691 .cloned()
692 .unwrap_or_default(),
693 endpoint: module.endpoint,
694 sink: module.sink.clone(),
695 negotiated_ver: module.negotiated_ver,
696 })
697 .collect())
698 }
699
700 pub(crate) fn request_connection_close(
704 &self,
705 connection_id: ConnectionId,
706 reason: CloseReason,
707 ) -> bool {
708 let sender = self.lock_close_registry().remove(&connection_id);
709 if let Some(sender) = sender {
710 debug!(
711 connection_id = connection_id.get(),
712 close_reason = %reason,
713 "requesting connection close"
714 );
715 let _ = sender.send(reason);
716 true
717 } else {
718 debug!(
719 connection_id = connection_id.get(),
720 close_reason = %reason,
721 "connection close request ignored for inactive connection"
722 );
723 false
724 }
725 }
726
727 pub fn register_module_connection(
728 &self,
729 connection_id: ConnectionId,
730 module_id: String,
731 negotiated_ver: u8,
732 concurrency: Concurrency,
733 sink: FrameSink,
734 ) -> Result<ModuleEndpointId, ForwardingError> {
735 self.register_module_connection_inner(
736 connection_id,
737 module_id,
738 negotiated_ver,
739 concurrency,
740 sink,
741 None,
742 )
743 }
744
745 pub(crate) fn register_module_connection_acked(
758 &self,
759 connection_id: ConnectionId,
760 module_id: String,
761 negotiated_ver: u8,
762 concurrency: Concurrency,
763 sink: FrameSink,
764 hello_ack: Frame,
765 ) -> Result<ModuleEndpointId, ForwardingError> {
766 self.register_module_connection_inner(
767 connection_id,
768 module_id,
769 negotiated_ver,
770 concurrency,
771 sink,
772 Some(hello_ack),
773 )
774 }
775
776 fn register_module_connection_inner(
777 &self,
778 connection_id: ConnectionId,
779 module_id: String,
780 negotiated_ver: u8,
781 concurrency: Concurrency,
782 sink: FrameSink,
783 hello_ack: Option<Frame>,
784 ) -> Result<ModuleEndpointId, ForwardingError> {
785 let mut inner = self.write_inner()?;
786 if inner.daemon_draining || inner.closing_connections.contains(&connection_id) {
787 return Err(ForwardingError::ConnectionClosing { connection_id });
788 }
789 check_module_connection_role_locked(&inner, connection_id)?;
790 enqueue_hello_ack_locked(&sink, connection_id, hello_ack)?;
793
794 inner.next_generation = inner.next_generation.checked_add(1).unwrap_or(1);
795 let endpoint = ModuleEndpointId {
796 connection_id,
797 generation: inner.next_generation,
798 };
799 inner.endpoint_by_connection.insert(connection_id, endpoint);
800 inner
801 .module_id_by_endpoint
802 .insert(endpoint, module_id.clone());
803 inner.next_module_channel.insert(endpoint, 1);
804 inner.next_control_corr.insert(endpoint, 1);
805 inner.modules_by_id.insert(
806 module_id.clone(),
807 ModuleConnection {
808 endpoint,
809 sink,
810 negotiated_ver,
811 concurrency,
812 },
813 );
814 drop(inner);
815
816 if let Some(discarded) = self
829 .route_bind_breakers
830 .reset_for_new_module_connection(&module_id)
831 {
832 info!(
833 module_id = %module_id,
834 discarded_consecutive_timeouts = discarded,
835 "route.bind breaker state discarded: a new module connection replaced the process it described"
836 );
837 }
838 Ok(endpoint)
839 }
840
841 #[cfg(test)]
855 pub(crate) fn register_candidate_module_connection(
856 &self,
857 connection_id: ConnectionId,
858 module_id: String,
859 negotiated_ver: u8,
860 concurrency: Concurrency,
861 sink: FrameSink,
862 ) -> Result<ModuleEndpointId, ForwardingError> {
863 self.register_candidate_module_connection_inner(
864 connection_id,
865 module_id,
866 negotiated_ver,
867 concurrency,
868 sink,
869 None,
870 )
871 }
872
873 pub(crate) fn register_candidate_module_connection_acked(
880 &self,
881 connection_id: ConnectionId,
882 module_id: String,
883 negotiated_ver: u8,
884 concurrency: Concurrency,
885 sink: FrameSink,
886 hello_ack: Frame,
887 ) -> Result<ModuleEndpointId, ForwardingError> {
888 self.register_candidate_module_connection_inner(
889 connection_id,
890 module_id,
891 negotiated_ver,
892 concurrency,
893 sink,
894 Some(hello_ack),
895 )
896 }
897
898 fn register_candidate_module_connection_inner(
899 &self,
900 connection_id: ConnectionId,
901 module_id: String,
902 negotiated_ver: u8,
903 concurrency: Concurrency,
904 sink: FrameSink,
905 hello_ack: Option<Frame>,
906 ) -> Result<ModuleEndpointId, ForwardingError> {
907 let mut inner = self.write_inner()?;
908 if inner.daemon_draining || inner.closing_connections.contains(&connection_id) {
909 return Err(ForwardingError::ConnectionClosing { connection_id });
910 }
911 check_module_connection_role_locked(&inner, connection_id)?;
912 if inner.candidates_by_id.contains_key(&module_id) {
913 return Err(ForwardingError::CandidateSlotOccupied { module_id });
914 }
915 enqueue_hello_ack_locked(&sink, connection_id, hello_ack)?;
918
919 inner.next_generation = inner.next_generation.checked_add(1).unwrap_or(1);
920 let endpoint = ModuleEndpointId {
921 connection_id,
922 generation: inner.next_generation,
923 };
924 inner.endpoint_by_connection.insert(connection_id, endpoint);
925 inner
926 .module_id_by_endpoint
927 .insert(endpoint, module_id.clone());
928 inner.next_module_channel.insert(endpoint, 1);
929 inner.next_control_corr.insert(endpoint, 1);
930 inner.candidates_by_id.insert(
931 module_id,
932 ModuleConnection {
933 endpoint,
934 sink,
935 negotiated_ver,
936 concurrency,
937 },
938 );
939 Ok(endpoint)
940 }
941
942 pub(crate) fn cutover_candidate(
960 &self,
961 module_id: &str,
962 ) -> Result<Option<ForwardingCutover>, ForwardingError> {
963 let mut inner = self.write_inner()?;
964 if inner.daemon_draining {
965 return Err(ForwardingError::ModuleReloading {
966 module_id: module_id.to_string(),
967 });
968 }
969 let Some(candidate) = inner.candidates_by_id.remove(module_id) else {
970 return Ok(None);
971 };
972 let promoted = candidate.endpoint;
973 let incumbent = inner.modules_by_id.insert(module_id.to_string(), candidate);
974 let incumbent = incumbent.map(|incumbent| {
975 let endpoint = incumbent.endpoint;
976 inner.superseded_endpoints.insert(endpoint, incumbent);
977 endpoint
978 });
979 drop(inner);
980
981 if let Some(discarded) = self
985 .route_bind_breakers
986 .reset_for_new_module_connection(module_id)
987 {
988 info!(
989 module_id = %module_id,
990 discarded_consecutive_timeouts = discarded,
991 "route.bind breaker state discarded: a swap candidate was promoted over the process it described"
992 );
993 }
994 Ok(Some(ForwardingCutover {
995 promoted,
996 incumbent,
997 }))
998 }
999
1000 #[allow(clippy::too_many_arguments)]
1001 pub(crate) async fn begin_route_bind_relay_for(
1002 &self,
1003 client_connection_id: ConnectionId,
1004 client_sink: FrameSink,
1005 client_negotiated_ver: u8,
1006 client_corr: u64,
1007 module_id: &str,
1008 principal: Principal,
1009 scope: Option<BoundScope>,
1010 project_root: Option<ProjectRootId>,
1011 deadline: Instant,
1012 ) -> Result<PendingRouteBindRelay, ForwardingError> {
1013 let client_permit =
1017 client_sink
1018 .reserve_owned()
1019 .await
1020 .map_err(|_| ForwardingError::ClientEgressClosed {
1021 connection_id: client_connection_id,
1022 })?;
1023 self.begin_route_bind_relay_inner(
1024 client_connection_id,
1025 client_sink,
1026 client_negotiated_ver,
1027 client_corr,
1028 module_id,
1029 principal,
1030 scope,
1031 project_root,
1032 deadline,
1033 client_permit,
1034 )
1035 }
1036
1037 #[cfg(test)]
1038 pub(crate) fn begin_route_bind_relay_for_test(
1039 &self,
1040 client_connection_id: ConnectionId,
1041 client_sink: FrameSink,
1042 client_corr: u64,
1043 module_id: &str,
1044 ) -> Result<PendingRouteBindRelay, ForwardingError> {
1045 let permit =
1046 client_sink
1047 .try_reserve_owned()
1048 .map_err(|_| ForwardingError::ClientEgressClosed {
1049 connection_id: client_connection_id,
1050 })?;
1051 self.begin_route_bind_relay_inner(
1052 client_connection_id,
1053 client_sink,
1054 subc_protocol::PROTOCOL_VERSION,
1055 client_corr,
1056 module_id,
1057 Principal::Direct,
1058 None,
1059 None,
1060 Instant::now() + std::time::Duration::from_secs(60),
1061 permit,
1062 )
1063 }
1064
1065 pub(crate) fn begin_module_control_rpc_for(
1066 &self,
1067 module_id: &str,
1068 expected_op: &str,
1069 deadline: Instant,
1070 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1071 self.begin_module_control_rpc_inner(module_id, expected_op, deadline, None, false)
1072 }
1073
1074 pub(crate) fn begin_health_probe_rpc_for(
1075 &self,
1076 module_id: &str,
1077 expected_op: &str,
1078 probe_started_at: Instant,
1079 deadline: Instant,
1080 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1081 self.begin_module_control_rpc_inner(
1082 module_id,
1083 expected_op,
1084 deadline,
1085 Some(probe_started_at),
1086 false,
1087 )
1088 }
1089
1090 pub(crate) fn begin_drain_health_probe_rpc_for(
1091 &self,
1092 module_id: &str,
1093 expected_op: &str,
1094 probe_started_at: Instant,
1095 deadline: Instant,
1096 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1097 self.begin_module_control_rpc_inner(
1098 module_id,
1099 expected_op,
1100 deadline,
1101 Some(probe_started_at),
1102 true,
1103 )
1104 }
1105
1106 pub(crate) fn begin_endpoint_health_probe_rpc_for(
1114 &self,
1115 endpoint: ModuleEndpointId,
1116 expected_op: &str,
1117 probe_started_at: Instant,
1118 deadline: Instant,
1119 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1120 let inner = self.write_inner()?;
1121 let module = module_connection_for_endpoint_locked(&inner, endpoint)
1122 .cloned()
1123 .ok_or(ForwardingError::NoModuleConnection)?;
1124 let module_id = inner
1125 .module_id_by_endpoint
1126 .get(&endpoint)
1127 .cloned()
1128 .unwrap_or_default();
1129 self.begin_control_rpc_locked(
1132 inner,
1133 &module_id,
1134 module,
1135 expected_op,
1136 deadline,
1137 Some(probe_started_at),
1138 true,
1139 )
1140 }
1141
1142 fn begin_module_control_rpc_inner(
1143 &self,
1144 module_id: &str,
1145 expected_op: &str,
1146 deadline: Instant,
1147 health_probe_started_at: Option<Instant>,
1148 allow_draining: bool,
1149 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1150 let inner = self.write_inner()?;
1151 let module = inner
1152 .modules_by_id
1153 .get(module_id)
1154 .cloned()
1155 .ok_or(ForwardingError::NoModuleConnection)?;
1156 self.begin_control_rpc_locked(
1157 inner,
1158 module_id,
1159 module,
1160 expected_op,
1161 deadline,
1162 health_probe_started_at,
1163 allow_draining,
1164 )
1165 }
1166
1167 #[allow(clippy::too_many_arguments)]
1168 fn begin_control_rpc_locked(
1169 &self,
1170 mut inner: RwLockWriteGuard<'_, ForwardingInner>,
1171 module_id: &str,
1172 module: ModuleConnection,
1173 expected_op: &str,
1174 deadline: Instant,
1175 health_probe_started_at: Option<Instant>,
1176 allow_draining: bool,
1177 ) -> Result<PendingModuleControlRpc, ForwardingError> {
1178 if !allow_draining && inner.draining_endpoints.contains_key(&module.endpoint) {
1179 return Err(ForwardingError::ModuleReloading {
1180 module_id: module_id.to_string(),
1181 });
1182 }
1183 if inner
1184 .closing_connections
1185 .contains(&module.endpoint.connection_id)
1186 {
1187 return Err(ForwardingError::ConnectionClosing {
1188 connection_id: module.endpoint.connection_id,
1189 });
1190 }
1191 if health_probe_started_at.is_some() {
1192 inner
1196 .health_probe_tombstones
1197 .retain(|(endpoint, _), _| *endpoint != module.endpoint);
1198 }
1199 let corr = match inner.allocate_control_corr(module.endpoint) {
1200 Ok(corr) => corr,
1201 Err(err) => {
1202 drop(inner);
1203 self.request_connection_close(
1204 module.endpoint.connection_id,
1205 CloseReason::new(
1206 "control_correlation_exhausted",
1207 "daemon-originated channel-0 correlation space exhausted",
1208 ),
1209 );
1210 return Err(err);
1211 }
1212 };
1213 let (sender, receiver) = oneshot::channel();
1214 inner.pending_control_rpcs.insert(
1215 (module.endpoint, corr),
1216 PendingModuleControlRpcEntry {
1217 expected_op: expected_op.to_string(),
1218 deadline,
1219 health_probe_started_at,
1220 sender,
1221 },
1222 );
1223
1224 Ok(PendingModuleControlRpc {
1225 endpoint: module.endpoint,
1226 module_sink: module.sink,
1227 negotiated_ver: module.negotiated_ver,
1228 corr,
1229 receiver,
1230 })
1231 }
1232
1233 #[allow(clippy::too_many_arguments)]
1234 fn begin_route_bind_relay_inner(
1235 &self,
1236 client_connection_id: ConnectionId,
1237 client_sink: FrameSink,
1238 client_negotiated_ver: u8,
1239 client_corr: u64,
1240 expected_module_id: &str,
1241 principal: Principal,
1242 scope: Option<BoundScope>,
1243 project_root: Option<ProjectRootId>,
1244 deadline: Instant,
1245 client_permit: crate::router::EgressPermit,
1246 ) -> Result<PendingRouteBindRelay, ForwardingError> {
1247 let mut inner = self.write_inner()?;
1248 if inner
1252 .endpoint_by_connection
1253 .contains_key(&client_connection_id)
1254 {
1255 return Err(ForwardingError::ConnectionRoleConflict {
1256 connection_id: client_connection_id,
1257 });
1258 }
1259 if inner.closing_connections.contains(&client_connection_id) {
1260 return Err(ForwardingError::ConnectionClosing {
1261 connection_id: client_connection_id,
1262 });
1263 }
1264 let module = inner
1265 .modules_by_id
1266 .get(expected_module_id)
1267 .cloned()
1268 .ok_or(ForwardingError::NoModuleConnection)?;
1269 if inner.draining_endpoints.contains_key(&module.endpoint) {
1270 return Err(ForwardingError::ModuleReloading {
1271 module_id: expected_module_id.to_string(),
1272 });
1273 }
1274 if inner
1275 .closing_connections
1276 .contains(&module.endpoint.connection_id)
1277 {
1278 return Err(ForwardingError::ConnectionClosing {
1279 connection_id: module.endpoint.connection_id,
1280 });
1281 }
1282
1283 let corr = match inner.allocate_control_corr(module.endpoint) {
1284 Ok(corr) => corr,
1285 Err(err) => {
1286 drop(inner);
1287 self.request_connection_close(
1288 module.endpoint.connection_id,
1289 CloseReason::new(
1290 "control_correlation_exhausted",
1291 "daemon-originated channel-0 correlation space exhausted",
1292 ),
1293 );
1294 return Err(err);
1295 }
1296 };
1297 let (client_channel, client_epoch, module_channel, module_epoch) =
1298 inner.allocate_route_slots(client_connection_id, module.endpoint)?;
1299 let client_key = ClientRouteKey {
1300 connection_id: client_connection_id,
1301 channel: client_channel,
1302 };
1303 let module_key = ModuleRouteKey {
1304 endpoint: module.endpoint,
1305 channel: module_channel,
1306 };
1307 let reservation = RouteReservation {
1308 client_key,
1309 module_key,
1310 client_epoch,
1311 module_epoch,
1312 project_root,
1313 };
1314 let response_body = serde_json::to_vec(&ClientControlResponse::RouteOpen {
1315 route_channel: client_channel,
1316 route_epoch: client_epoch,
1317 })
1318 .map_err(|err| ForwardingError::RouteOpenBuild(err.to_string()))?;
1319 let route_open_frame = Frame::build_with_version(
1320 client_negotiated_ver,
1321 FrameType::Response,
1322 Flags::new(false, Priority::Passive, false),
1323 0,
1324 0,
1325 client_corr,
1326 response_body,
1327 )
1328 .map_err(|err| ForwardingError::RouteOpenBuild(err.to_string()))?;
1329 let (sender, receiver) = oneshot::channel();
1330 inner.reserved_client.insert(client_key, module_key);
1331 inner.reserved_module.insert(module_key, client_key);
1332 inner.pending_relays.insert(
1333 (module.endpoint, corr),
1334 PendingRouteBindRelayEntry {
1335 reservation,
1336 client_sink,
1337 client_negotiated_ver,
1338 client_permit,
1339 route_open_frame,
1340 principal,
1341 scope,
1342 deadline,
1343 relay_enqueued: false,
1344 sender,
1345 },
1346 );
1347
1348 Ok(PendingRouteBindRelay {
1349 endpoint: module.endpoint,
1350 module_sink: module.sink,
1351 negotiated_ver: module.negotiated_ver,
1352 client_channel,
1353 client_epoch,
1354 module_channel,
1355 module_epoch,
1356 corr,
1357 receiver,
1358 })
1359 }
1360
1361 pub(crate) fn mark_route_bind_relay_enqueued(
1362 &self,
1363 endpoint: ModuleEndpointId,
1364 corr: u64,
1365 ) -> Result<bool, ForwardingError> {
1366 let mut inner = self.write_inner()?;
1367 let Some(pending) = inner.pending_relays.get_mut(&(endpoint, corr)) else {
1368 return Ok(false);
1369 };
1370 pending.relay_enqueued = true;
1371 Ok(true)
1372 }
1373
1374 pub(crate) fn release_client_route(
1375 &self,
1376 client_connection_id: ConnectionId,
1377 client_channel: u16,
1378 expected_epoch: u32,
1379 ) -> Result<RouteRelease, ForwardingError> {
1380 let mut inner = self.write_inner()?;
1381 let release = release_client_route_locked(
1382 &mut inner,
1383 ClientRouteKey {
1384 connection_id: client_connection_id,
1385 channel: client_channel,
1386 },
1387 expected_epoch,
1388 );
1389 self.record_route_release(&release);
1390 Ok(release)
1391 }
1392
1393 pub(crate) fn release_module_route(
1394 &self,
1395 module_connection_id: ConnectionId,
1396 module_channel: u16,
1397 expected_epoch: u32,
1398 ) -> Result<RouteRelease, ForwardingError> {
1399 let mut inner = self.write_inner()?;
1400 let Some(endpoint) = inner
1401 .endpoint_by_connection
1402 .get(&module_connection_id)
1403 .copied()
1404 else {
1405 return Ok(RouteRelease::Absent);
1406 };
1407 let release = release_module_route_locked(
1408 &mut inner,
1409 ModuleRouteKey {
1410 endpoint,
1411 channel: module_channel,
1412 },
1413 expected_epoch,
1414 );
1415 self.record_route_release(&release);
1416 Ok(release)
1417 }
1418
1419 pub(crate) fn abort_pending_relay(
1420 &self,
1421 endpoint: ModuleEndpointId,
1422 corr: u64,
1423 outcome: RouteBindRelayOutcome,
1424 ) -> Result<Option<GoodbyeTarget>, ForwardingError> {
1425 let mut inner = self.write_inner()?;
1426 let Some(pending) = inner.pending_relays.remove(&(endpoint, corr)) else {
1427 return Ok(None);
1428 };
1429 release_reserved_route_locked(
1430 &mut inner,
1431 pending.reservation.client_key,
1432 pending.reservation.module_key,
1433 );
1434 let target = pending
1435 .relay_enqueued
1436 .then(|| abandoned_route_target(&inner, &pending.reservation));
1437 let _ = pending.sender.send(outcome);
1438 Ok(target.flatten())
1439 }
1440
1441 pub(crate) fn cancel_module_control_rpc(
1442 &self,
1443 endpoint: ModuleEndpointId,
1444 corr: u64,
1445 ) -> Result<(), ForwardingError> {
1446 self.write_inner()?
1447 .pending_control_rpcs
1448 .remove(&(endpoint, corr));
1449 Ok(())
1450 }
1451
1452 pub(crate) fn tombstone_health_probe_rpc(
1453 &self,
1454 endpoint: ModuleEndpointId,
1455 corr: u64,
1456 ) -> Result<bool, ForwardingError> {
1457 let key = (endpoint, corr);
1458 let expires_at = Instant::now() + HEALTH_PROBE_TOMBSTONE_TTL;
1459 {
1460 let mut inner = self.write_inner()?;
1461 let Some(pending) = inner.pending_control_rpcs.remove(&key) else {
1462 return Ok(false);
1463 };
1464 let Some(probe_started_at) = pending.health_probe_started_at else {
1465 inner.pending_control_rpcs.insert(key, pending);
1466 return Ok(false);
1467 };
1468 let module_id = inner
1469 .module_id_by_endpoint
1470 .get(&endpoint)
1471 .cloned()
1472 .unwrap_or_else(|| "unknown".to_string());
1473 inner.health_probe_tombstones.insert(
1474 key,
1475 HealthProbeTombstone {
1476 expected_op: pending.expected_op,
1477 module_id,
1478 probe_started_at,
1479 expires_at,
1480 },
1481 );
1482 }
1483 self.schedule_health_probe_tombstone_expiration(key, expires_at);
1484 Ok(true)
1485 }
1486
1487 fn schedule_health_probe_tombstone_expiration(
1488 &self,
1489 key: (ModuleEndpointId, u64),
1490 expires_at: Instant,
1491 ) {
1492 let inner = Arc::downgrade(&self.inner);
1493 tokio::spawn(async move {
1494 tokio::time::sleep_until(expires_at).await;
1495 let Some(inner) = inner.upgrade() else {
1496 return;
1497 };
1498 let Ok(mut inner) = inner.write() else {
1499 return;
1500 };
1501 let expired = inner
1502 .health_probe_tombstones
1503 .get(&key)
1504 .is_some_and(|tombstone| tombstone.expires_at <= Instant::now());
1505 if expired {
1506 inner.health_probe_tombstones.remove(&key);
1507 }
1508 });
1509 }
1510
1511 pub(crate) fn complete_pending_relay(
1512 &self,
1513 connection_id: ConnectionId,
1514 corr: u64,
1515 outcome: RouteBindRelayOutcome,
1516 ) -> Result<PendingRelayCompletion, ForwardingError> {
1517 let mut inner = self.write_inner()?;
1518 let Some(endpoint) = inner.endpoint_by_connection.get(&connection_id).copied() else {
1519 return Ok(PendingRelayCompletion {
1520 settled: false,
1521 abandoned: None,
1522 });
1523 };
1524 let Some(pending) = inner.pending_relays.remove(&(endpoint, corr)) else {
1525 return Ok(PendingRelayCompletion {
1526 settled: false,
1527 abandoned: None,
1528 });
1529 };
1530
1531 if Instant::now() >= pending.deadline {
1532 release_reserved_route_locked(
1533 &mut inner,
1534 pending.reservation.client_key,
1535 pending.reservation.module_key,
1536 );
1537 let abandoned = matches!(outcome, RouteBindRelayOutcome::Accepted)
1538 .then(|| abandoned_route_target(&inner, &pending.reservation))
1539 .flatten();
1540 let _ = pending
1541 .sender
1542 .send(RouteBindRelayOutcome::Rejected(ErrorBody {
1543 code: "module_timeout".to_string(),
1544 message: "route.bind response arrived after its daemon deadline".to_string(),
1545 detail: None,
1546 }));
1547 return Ok(PendingRelayCompletion {
1548 settled: true,
1549 abandoned,
1550 });
1551 }
1552
1553 match outcome {
1554 RouteBindRelayOutcome::Accepted
1573 if pending.client_sink.is_closed()
1574 || inner
1575 .closing_connections
1576 .contains(&pending.reservation.client_key.connection_id) =>
1577 {
1578 let reason = if pending.client_sink.is_closed() {
1579 "client egress closed before route publication"
1580 } else {
1581 "client connection is closing before route publication"
1582 };
1583 release_reserved_route_locked(
1584 &mut inner,
1585 pending.reservation.client_key,
1586 pending.reservation.module_key,
1587 );
1588 let abandoned = pending
1589 .relay_enqueued
1590 .then(|| abandoned_route_target(&inner, &pending.reservation))
1591 .flatten();
1592 let _ = pending
1593 .sender
1594 .send(RouteBindRelayOutcome::ModuleGone(reason.to_string()));
1595 return Ok(PendingRelayCompletion {
1596 settled: true,
1597 abandoned,
1598 });
1599 }
1600 RouteBindRelayOutcome::Accepted
1623 if inner.superseded_endpoints.contains_key(&endpoint)
1624 || inner.draining_endpoints.contains_key(&endpoint) =>
1625 {
1626 release_reserved_route_locked(
1627 &mut inner,
1628 pending.reservation.client_key,
1629 pending.reservation.module_key,
1630 );
1631 let abandoned = abandoned_route_target(&inner, &pending.reservation);
1634 let module_id = inner
1635 .module_id_by_endpoint
1636 .get(&endpoint)
1637 .cloned()
1638 .unwrap_or_else(|| "unknown".to_string());
1639 let _ = pending
1640 .sender
1641 .send(RouteBindRelayOutcome::Rejected(ErrorBody::new(
1642 "module_reloading",
1643 format!("module_id '{module_id}' is reloading"),
1644 )));
1645 return Ok(PendingRelayCompletion {
1646 settled: true,
1647 abandoned,
1648 });
1649 }
1650 RouteBindRelayOutcome::Accepted
1659 if pending
1660 .scope
1661 .as_ref()
1662 .is_some_and(|scope| scope_refusal_locked(&inner, scope).is_some()) =>
1663 {
1664 let refusal = pending
1665 .scope
1666 .as_ref()
1667 .and_then(|scope| scope_refusal_locked(&inner, scope))
1668 .expect("guard matched a refusal under the same lock");
1669 release_reserved_route_locked(
1670 &mut inner,
1671 pending.reservation.client_key,
1672 pending.reservation.module_key,
1673 );
1674 let abandoned = abandoned_route_target(&inner, &pending.reservation);
1675 let _ = pending
1676 .sender
1677 .send(RouteBindRelayOutcome::Rejected(refusal));
1678 return Ok(PendingRelayCompletion {
1679 settled: true,
1680 abandoned,
1681 });
1682 }
1683 RouteBindRelayOutcome::Accepted => {
1684 let abandoned = commit_route_locked(&mut inner, pending)?;
1685 return Ok(PendingRelayCompletion {
1686 settled: true,
1687 abandoned,
1688 });
1689 }
1690 terminal => {
1691 release_reserved_route_locked(
1692 &mut inner,
1693 pending.reservation.client_key,
1694 pending.reservation.module_key,
1695 );
1696 let _ = pending.sender.send(terminal);
1697 }
1698 }
1699 Ok(PendingRelayCompletion {
1700 settled: true,
1701 abandoned: None,
1702 })
1703 }
1704
1705 pub(crate) fn pending_module_control_op(
1706 &self,
1707 connection_id: ConnectionId,
1708 corr: u64,
1709 ) -> Result<Option<String>, ForwardingError> {
1710 let inner = self.read_inner()?;
1711 let Some(endpoint) = inner.endpoint_by_connection.get(&connection_id).copied() else {
1712 return Ok(None);
1713 };
1714 let key = (endpoint, corr);
1715 Ok(inner
1716 .pending_control_rpcs
1717 .get(&key)
1718 .map(|pending| pending.expected_op.clone())
1719 .or_else(|| {
1720 inner
1721 .health_probe_tombstones
1722 .get(&key)
1723 .filter(|tombstone| tombstone.expires_at > Instant::now())
1724 .map(|tombstone| tombstone.expected_op.clone())
1725 }))
1726 }
1727
1728 pub(crate) fn complete_module_control_rpc(
1729 &self,
1730 connection_id: ConnectionId,
1731 corr: u64,
1732 actual_op: Option<&str>,
1733 outcome: ModuleControlRpcOutcome,
1734 ) -> Result<ModuleControlRpcCompletion, ForwardingError> {
1735 let now = Instant::now();
1736 let mut inner = self.write_inner()?;
1737 let Some(endpoint) = inner.endpoint_by_connection.get(&connection_id).copied() else {
1738 return Ok(ModuleControlRpcCompletion::Unknown);
1739 };
1740 let key = (endpoint, corr);
1741 if let Some(pending) = inner.pending_control_rpcs.remove(&key) {
1742 if now >= pending.deadline {
1743 let late_health_answer = pending.health_probe_started_at.map(|probe_started_at| {
1744 ModuleControlRpcCompletion::LateHealthAnswer {
1745 module_id: inner
1746 .module_id_by_endpoint
1747 .get(&endpoint)
1748 .cloned()
1749 .unwrap_or_else(|| "unknown".to_string()),
1750 latency: now.saturating_duration_since(probe_started_at),
1751 }
1752 });
1753 let _ = pending
1754 .sender
1755 .send(ModuleControlRpcOutcome::DeadlineElapsed);
1756 return Ok(late_health_answer.unwrap_or(ModuleControlRpcCompletion::Settled));
1757 }
1758 let outcome = match actual_op {
1759 Some(actual) if actual != pending.expected_op => {
1760 ModuleControlRpcOutcome::UnexpectedOp {
1761 expected: pending.expected_op,
1762 actual: actual.to_string(),
1763 }
1764 }
1765 _ => outcome,
1766 };
1767 let _ = pending.sender.send(outcome);
1768 return Ok(ModuleControlRpcCompletion::Settled);
1769 }
1770
1771 let Some(tombstone) = inner.health_probe_tombstones.remove(&key) else {
1772 return Ok(ModuleControlRpcCompletion::Unknown);
1773 };
1774 if tombstone.expires_at <= now {
1775 return Ok(ModuleControlRpcCompletion::Unknown);
1776 }
1777 Ok(ModuleControlRpcCompletion::LateHealthAnswer {
1778 module_id: tombstone.module_id,
1779 latency: now.saturating_duration_since(tombstone.probe_started_at),
1780 })
1781 }
1782
1783 #[cfg(test)]
1784 pub(crate) fn health_probe_tombstone_count(&self) -> Result<usize, ForwardingError> {
1785 Ok(self.read_inner()?.health_probe_tombstones.len())
1786 }
1787
1788 #[cfg(test)]
1789 pub(crate) fn closing_connection_count(&self) -> Result<usize, ForwardingError> {
1790 Ok(self.read_inner()?.closing_connections.len())
1791 }
1792
1793 #[cfg(test)]
1796 pub(crate) fn reserved_route_count(&self) -> Result<(usize, usize), ForwardingError> {
1797 let inner = self.read_inner()?;
1798 Ok((inner.reserved_client.len(), inner.reserved_module.len()))
1799 }
1800
1801 pub(crate) fn module_endpoint_for_connection(
1802 &self,
1803 connection_id: ConnectionId,
1804 ) -> Result<Option<ModuleEndpointId>, ForwardingError> {
1805 Ok(self
1806 .read_inner()?
1807 .endpoint_by_connection
1808 .get(&connection_id)
1809 .copied())
1810 }
1811
1812 pub(crate) fn module_id_for_connection(
1815 &self,
1816 connection_id: ConnectionId,
1817 ) -> Result<Option<String>, ForwardingError> {
1818 let inner = self.read_inner()?;
1819 Ok(inner
1820 .endpoint_by_connection
1821 .get(&connection_id)
1822 .and_then(|endpoint| inner.module_id_by_endpoint.get(endpoint))
1823 .cloned())
1824 }
1825
1826 pub(crate) fn module_route_epoch_was_allocated(
1833 &self,
1834 connection_id: ConnectionId,
1835 channel: u16,
1836 epoch: u32,
1837 ) -> Result<bool, ForwardingError> {
1838 let inner = self.read_inner()?;
1839 let Some(endpoint) = inner.endpoint_by_connection.get(&connection_id).copied() else {
1840 return Ok(false);
1841 };
1842 Ok(inner
1843 .module_slot_epochs
1844 .get(&ModuleRouteKey { endpoint, channel })
1845 .is_some_and(|last| epoch != 0 && epoch <= *last))
1846 }
1847
1848 pub(crate) fn has_live_module_connection(
1849 &self,
1850 module_id: &str,
1851 ) -> Result<bool, ForwardingError> {
1852 Ok(self.read_inner()?.modules_by_id.contains_key(module_id))
1853 }
1854
1855 pub(crate) fn lookup_data_route(
1856 &self,
1857 connection_id: ConnectionId,
1858 channel: u16,
1859 epoch: u32,
1860 ) -> Result<DataRoute, ForwardingError> {
1861 let inner = self.read_inner()?;
1862 let state = if let Some(endpoint) =
1863 inner.endpoint_by_connection.get(&connection_id).copied()
1864 {
1865 let key = ModuleRouteKey { endpoint, channel };
1866 match inner.module_to_client.get(&key) {
1867 Some(route) if route.module_epoch == epoch => {
1868 DataRouteState::Bound(Arc::clone(route))
1869 }
1870 Some(_) => DataRouteState::EpochMismatch,
1871 None if inner.reserved_module.contains_key(&key)
1872 && inner.module_slot_epochs.get(&key).copied() == Some(epoch) =>
1873 {
1874 DataRouteState::Reserved
1875 }
1876 None if inner.reserved_module.contains_key(&key) => DataRouteState::EpochMismatch,
1877 None => DataRouteState::Absent,
1878 }
1879 } else {
1880 let key = ClientRouteKey {
1881 connection_id,
1882 channel,
1883 };
1884 match inner.client_to_module.get(&key) {
1885 Some(route) if route.client_epoch == epoch => {
1886 DataRouteState::Bound(Arc::clone(route))
1887 }
1888 Some(_) => DataRouteState::EpochMismatch,
1889 None if inner.reserved_client.contains_key(&key)
1890 && inner.client_slot_epochs.get(&key).copied() == Some(epoch) =>
1891 {
1892 DataRouteState::Reserved
1893 }
1894 None if inner.reserved_client.contains_key(&key) => DataRouteState::EpochMismatch,
1895 None => DataRouteState::Absent,
1896 }
1897 };
1898 Ok(
1899 if inner.endpoint_by_connection.contains_key(&connection_id) {
1900 DataRoute::Module(state)
1901 } else {
1902 DataRoute::Client(state)
1903 },
1904 )
1905 }
1906
1907 #[cfg(test)]
1908 pub(crate) fn inject_client_slot_epoch(
1909 &self,
1910 connection_id: ConnectionId,
1911 channel: u16,
1912 last_epoch: u32,
1913 ) {
1914 let mut inner = self.write_inner().expect("forwarding lock");
1915 inner.client_slot_epochs.insert(
1916 ClientRouteKey {
1917 connection_id,
1918 channel,
1919 },
1920 last_epoch,
1921 );
1922 inner.next_client_channel.insert(connection_id, channel);
1923 }
1924
1925 #[cfg(test)]
1926 pub(crate) fn inject_module_slot_epoch(
1927 &self,
1928 endpoint: ModuleEndpointId,
1929 channel: u16,
1930 last_epoch: u32,
1931 ) {
1932 let mut inner = self.write_inner().expect("forwarding lock");
1933 inner
1934 .module_slot_epochs
1935 .insert(ModuleRouteKey { endpoint, channel }, last_epoch);
1936 inner.next_module_channel.insert(endpoint, channel);
1937 }
1938
1939 #[cfg(test)]
1940 pub(crate) fn inject_control_corr(&self, endpoint: ModuleEndpointId, next_corr: u64) {
1941 self.write_inner()
1942 .expect("forwarding lock")
1943 .next_control_corr
1944 .insert(endpoint, next_corr);
1945 }
1946
1947 pub(crate) fn cache_status(
1948 &self,
1949 endpoint: ModuleEndpointId,
1950 module_channel: u16,
1951 module_epoch: u32,
1952 status: String,
1953 ) -> Result<bool, ForwardingError> {
1954 let mut inner = self.write_inner()?;
1955 if !inner.module_id_by_endpoint.contains_key(&endpoint) {
1956 return Err(ForwardingError::StaleModuleEndpoint);
1957 }
1958
1959 let module_key = ModuleRouteKey {
1960 endpoint,
1961 channel: module_channel,
1962 };
1963 let handle = if let Some(route) = inner.module_to_client.get(&module_key) {
1964 (route.module_epoch == module_epoch).then_some((
1965 ClientRouteKey {
1966 connection_id: route.client_connection_id,
1967 channel: route.client_channel,
1968 },
1969 route.client_epoch,
1970 ))
1971 } else if let Some(client_key) = inner.reserved_module.get(&module_key).copied() {
1972 (inner.module_slot_epochs.get(&module_key).copied() == Some(module_epoch)).then_some((
1973 client_key,
1974 inner
1975 .client_slot_epochs
1976 .get(&client_key)
1977 .copied()
1978 .unwrap_or(0),
1979 ))
1980 } else {
1981 None
1982 };
1983
1984 if let Some(handle) = handle {
1985 inner.status.insert(handle, status);
1986 Ok(true)
1987 } else {
1988 debug!(
1989 module_channel,
1990 module_epoch,
1991 generation = endpoint.generation,
1992 connection_id = endpoint.connection_id.get(),
1993 "dropping stale status update for module route handle"
1994 );
1995 Ok(false)
1996 }
1997 }
1998
1999 pub(crate) fn route_poll_snapshot(
2000 &self,
2001 client_connection_id: ConnectionId,
2002 client_channel: u16,
2003 client_epoch: u32,
2004 ) -> Result<RoutePollSnapshot, ForwardingError> {
2005 let inner = self.read_inner()?;
2006 let client_key = ClientRouteKey {
2007 connection_id: client_connection_id,
2008 channel: client_channel,
2009 };
2010 let Some(route) = inner.client_to_module.get(&client_key) else {
2011 return Ok(RoutePollSnapshot::Absent);
2012 };
2013 if route.client_epoch != client_epoch
2014 || !inner
2015 .module_id_by_endpoint
2016 .contains_key(&route.module_endpoint)
2017 {
2018 return Ok(RoutePollSnapshot::Absent);
2019 }
2020 Ok(RoutePollSnapshot::Bound {
2021 module_id: route.module_id.clone(),
2022 status: inner.status.get(&(client_key, client_epoch)).cloned(),
2023 })
2024 }
2025
2026 pub fn active_binding_count(&self) -> Result<usize, ForwardingError> {
2027 Ok(self.read_inner()?.client_to_module.len())
2028 }
2029
2030 pub fn client_route_concentration(&self) -> Result<(usize, usize), ForwardingError> {
2040 let inner = self.read_inner()?;
2041 let mut per_connection: HashMap<ConnectionId, usize> = HashMap::new();
2042 for key in inner.client_to_module.keys() {
2043 *per_connection.entry(key.connection_id).or_insert(0) += 1;
2044 }
2045 let max = per_connection.values().copied().max().unwrap_or(0);
2046 Ok((per_connection.len(), max))
2047 }
2048
2049 pub fn has_route_channel(&self, route_channel: u16) -> Result<bool, ForwardingError> {
2050 let inner = self.read_inner()?;
2051 Ok(inner
2052 .client_to_module
2053 .keys()
2054 .any(|key| key.channel == route_channel))
2055 }
2056
2057 pub(crate) fn is_daemon_draining(&self) -> Result<bool, ForwardingError> {
2059 Ok(self.read_inner()?.daemon_draining)
2060 }
2061
2062 #[cfg(unix)]
2065 pub(crate) fn begin_daemon_drain(&self) -> Result<Vec<String>, ForwardingError> {
2066 let mut inner = self.write_inner()?;
2067 inner.daemon_draining = true;
2068 let modules = inner
2069 .modules_by_id
2070 .iter()
2071 .map(|(id, module)| (id.clone(), module.endpoint))
2072 .collect::<Vec<_>>();
2073 for (_, endpoint) in &modules {
2074 inner
2075 .draining_endpoints
2076 .insert(*endpoint, RouteCloseReason::Restart);
2077 }
2078 let off_slot_endpoints = inner
2083 .candidates_by_id
2084 .values()
2085 .map(|module| module.endpoint)
2086 .chain(inner.superseded_endpoints.keys().copied())
2087 .collect::<Vec<_>>();
2088 for endpoint in off_slot_endpoints {
2089 inner
2090 .draining_endpoints
2091 .insert(endpoint, RouteCloseReason::Restart);
2092 }
2093 Ok(modules.into_iter().map(|(id, _)| id).collect())
2094 }
2095
2096 pub(crate) fn begin_module_drain(
2103 &self,
2104 module_id: &str,
2105 reason: RouteCloseReason,
2106 ) -> Result<Option<ModuleDrainTarget>, ForwardingError> {
2107 let mut inner = self.write_inner()?;
2108 let Some(module) = inner.modules_by_id.get(module_id).cloned() else {
2109 return Ok(None);
2110 };
2111 Ok(Some(begin_drain_locked(
2112 &mut inner, module_id, module, reason,
2113 )))
2114 }
2115
2116 pub(crate) fn begin_endpoint_drain(
2123 &self,
2124 endpoint: ModuleEndpointId,
2125 reason: RouteCloseReason,
2126 ) -> Result<Option<ModuleDrainTarget>, ForwardingError> {
2127 let mut inner = self.write_inner()?;
2128 let Some(module) = module_connection_for_endpoint_locked(&inner, endpoint).cloned() else {
2129 return Ok(None);
2130 };
2131 let module_id = inner
2132 .module_id_by_endpoint
2133 .get(&endpoint)
2134 .cloned()
2135 .expect("an endpoint resolved to a module connection has a module id");
2136 Ok(Some(begin_drain_locked(
2137 &mut inner, &module_id, module, reason,
2138 )))
2139 }
2140}
2141
2142fn begin_drain_locked(
2146 inner: &mut ForwardingInner,
2147 module_id: &str,
2148 module: ModuleConnection,
2149 reason: RouteCloseReason,
2150) -> ModuleDrainTarget {
2151 {
2152 let endpoint = module.endpoint;
2153 inner.draining_endpoints.insert(endpoint, reason);
2154
2155 let flows = inner
2156 .client_to_module
2157 .values()
2158 .filter(|route| route.module_endpoint == endpoint)
2159 .map(|route| Arc::clone(&route.flow))
2160 .collect::<Vec<_>>();
2161 let excluded_subscriptions = flows
2162 .into_iter()
2163 .map(|flow| flow.begin_drain())
2164 .fold(0u32, u32::saturating_add);
2165
2166 let pending_keys = inner
2167 .pending_relays
2168 .keys()
2169 .filter(|(pending_endpoint, _)| *pending_endpoint == endpoint)
2170 .copied()
2171 .collect::<Vec<_>>();
2172 let mut abandoned_bindings = Vec::new();
2173 for key in pending_keys {
2174 let Some(pending) = inner.pending_relays.remove(&key) else {
2175 continue;
2176 };
2177 release_reserved_route_locked(
2178 inner,
2179 pending.reservation.client_key,
2180 pending.reservation.module_key,
2181 );
2182 if pending.relay_enqueued {
2183 if let Some(target) = abandoned_route_target(inner, &pending.reservation) {
2184 abandoned_bindings.push(target);
2185 }
2186 }
2187 let _ = pending
2188 .sender
2189 .send(RouteBindRelayOutcome::Rejected(ErrorBody::new(
2190 "module_reloading",
2191 format!("module_id '{module_id}' is reloading"),
2192 )));
2193 }
2194
2195 let pending_control_keys = inner
2196 .pending_control_rpcs
2197 .keys()
2198 .filter(|(pending_endpoint, _)| *pending_endpoint == endpoint)
2199 .copied()
2200 .collect::<Vec<_>>();
2201 for key in pending_control_keys {
2202 if let Some(pending) = inner.pending_control_rpcs.remove(&key) {
2203 let _ = pending
2204 .sender
2205 .send(ModuleControlRpcOutcome::ModuleGone(format!(
2206 "module '{module_id}' began draining during module-control RPC"
2207 )));
2208 }
2209 }
2210
2211 ModuleDrainTarget {
2212 endpoint,
2213 sink: module.sink,
2214 negotiated_ver: module.negotiated_ver,
2215 abandoned_bindings,
2216 excluded_subscriptions,
2217 }
2218 }
2219}
2220
2221#[derive(Debug, Default, PartialEq, Eq)]
2224pub(crate) struct DrainHoldouts {
2225 pub(crate) requests: usize,
2228 pub(crate) routes: usize,
2230 pub(crate) total_routes: usize,
2232 pub(crate) top_connections: Vec<(u64, usize)>,
2235 pub(crate) held: Vec<(u16, u64)>,
2244}
2245
2246pub(crate) const DRAIN_HELD_REQUESTS_LISTED: usize = 32;
2248
2249impl ForwardingTable {
2250 pub(crate) fn endpoint_drain_holdouts(
2252 &self,
2253 endpoint: ModuleEndpointId,
2254 ) -> Result<DrainHoldouts, ForwardingError> {
2255 let inner = self.read_inner()?;
2256 let mut holdouts = DrainHoldouts::default();
2257 let mut by_connection: HashMap<u64, usize> = HashMap::new();
2258 for (key, route) in &inner.client_to_module {
2259 if route.module_endpoint != endpoint {
2260 continue;
2261 }
2262 holdouts.total_routes += 1;
2263 let held = route.flow.drain_in_flight();
2264 if held == 0 {
2265 continue;
2266 }
2267 holdouts.requests += held;
2268 holdouts.routes += 1;
2269 *by_connection.entry(key.connection_id.get()).or_default() += held;
2270 holdouts.held.extend(
2271 route
2272 .flow
2273 .drain_held_corrs()
2274 .into_iter()
2275 .map(|corr| (route.module_channel, corr)),
2276 );
2277 }
2278 holdouts.held.sort_unstable();
2279 holdouts.held.truncate(DRAIN_HELD_REQUESTS_LISTED);
2280 let mut connections = by_connection.into_iter().collect::<Vec<_>>();
2281 connections.sort_by(|left, right| right.1.cmp(&left.1).then(left.0.cmp(&right.0)));
2282 connections.truncate(3);
2283 holdouts.top_connections = connections;
2284 Ok(holdouts)
2285 }
2286
2287 pub(crate) fn endpoint_in_flight_count(
2288 &self,
2289 endpoint: ModuleEndpointId,
2290 ) -> Result<usize, ForwardingError> {
2291 let inner = self.read_inner()?;
2292 Ok(inner
2293 .client_to_module
2294 .values()
2295 .filter(|route| route.module_endpoint == endpoint)
2296 .map(|route| route.flow.drain_in_flight())
2297 .sum())
2298 }
2299
2300 pub(crate) fn endpoint_is_draining(
2301 &self,
2302 endpoint: ModuleEndpointId,
2303 ) -> Result<bool, ForwardingError> {
2304 Ok(self
2305 .read_inner()?
2306 .draining_endpoints
2307 .contains_key(&endpoint))
2308 }
2309
2310 pub(crate) fn module_is_draining(&self, module_id: &str) -> Result<bool, ForwardingError> {
2311 let inner = self.read_inner()?;
2312 Ok(inner
2313 .modules_by_id
2314 .get(module_id)
2315 .is_some_and(|module| inner.draining_endpoints.contains_key(&module.endpoint)))
2316 }
2317
2318 pub(crate) fn release_module_endpoint_routes(
2319 &self,
2320 endpoint: ModuleEndpointId,
2321 ) -> Result<Vec<GoodbyeTarget>, ForwardingError> {
2322 let mut inner = self.write_inner()?;
2323 let routes = inner
2324 .module_to_client
2325 .iter()
2326 .filter(|(module_key, _)| module_key.endpoint == endpoint)
2327 .map(|(module_key, route)| (*module_key, route.module_epoch))
2328 .collect::<Vec<_>>();
2329 let mut released = Vec::with_capacity(routes.len());
2330 for (module_key, epoch) in routes {
2331 if let RouteRelease::Removed(target) =
2332 release_module_route_locked(&mut inner, module_key, epoch)
2333 {
2334 released.push(target);
2335 }
2336 }
2337 Ok(released)
2338 }
2339
2340 pub(crate) fn endpoint_routes(
2346 &self,
2347 endpoint: ModuleEndpointId,
2348 ) -> Result<Vec<EndpointRoute>, ForwardingError> {
2349 let inner = self.read_inner()?;
2350 Ok(endpoint_routes_locked(&inner, endpoint))
2351 }
2352
2353 pub(crate) fn route_census(
2355 &self,
2356 module_id: Option<&str>,
2357 ) -> Result<Vec<(String, Vec<EndpointRoute>)>, ForwardingError> {
2358 let inner = self.read_inner()?;
2359 let mut endpoints = inner
2360 .modules_by_id
2361 .iter()
2362 .filter(|(id, _)| module_id.is_none_or(|requested| requested == id.as_str()))
2363 .map(|(id, module)| (id.clone(), module.endpoint))
2364 .collect::<Vec<_>>();
2365 endpoints.sort_by(|left, right| left.0.cmp(&right.0));
2366 Ok(endpoints
2367 .into_iter()
2368 .map(|(id, endpoint)| (id, endpoint_routes_locked(&inner, endpoint)))
2369 .collect())
2370 }
2371
2372 pub(crate) fn live_roots(
2374 &self,
2375 module_id: &str,
2376 ) -> Result<ModuleControlResponseToModule, ForwardingError> {
2377 let inner = self.read_inner()?;
2378 let endpoint = inner
2379 .modules_by_id
2380 .get(module_id)
2381 .map(|module| module.endpoint);
2382 let mut roots = BTreeMap::new();
2383 let mut unknown_root_bindings = 0;
2384 let mut total_bindings = 0;
2385 if let Some(endpoint) = endpoint {
2386 for binding in inner
2387 .module_to_client
2388 .values()
2389 .filter(|binding| binding.module_endpoint == endpoint)
2390 {
2391 total_bindings += 1;
2392 if let Some(root) = &binding.project_root {
2393 let entry = roots.entry(root.as_path().to_path_buf()).or_insert((0, 0));
2394 entry.0 += 1;
2395 } else {
2396 unknown_root_bindings += 1;
2397 }
2398 }
2399 for pending in inner
2400 .pending_relays
2401 .values()
2402 .filter(|pending| pending.reservation.module_key.endpoint == endpoint)
2403 {
2404 total_bindings += 1;
2405 if let Some(root) = &pending.reservation.project_root {
2406 let entry = roots.entry(root.as_path().to_path_buf()).or_insert((0, 0));
2407 entry.1 += 1;
2408 } else {
2409 unknown_root_bindings += 1;
2410 }
2411 }
2412 }
2413 Ok(ModuleControlResponseToModule::LiveRoots {
2414 roots: roots
2415 .into_iter()
2416 .map(|(project_root, (bound, pending))| LiveRoot {
2417 project_root,
2418 bound,
2419 pending,
2420 })
2421 .collect(),
2422 unknown_root_bindings,
2423 total_bindings,
2424 })
2425 }
2426
2427 pub(crate) fn connection_has_client_routes(
2433 &self,
2434 connection_id: ConnectionId,
2435 ) -> Result<bool, ForwardingError> {
2436 let inner = self.read_inner()?;
2437 Ok(connection_has_client_routes_locked(&inner, connection_id))
2438 }
2439
2440 pub(crate) fn cleanup_connection(
2441 &self,
2442 connection_id: ConnectionId,
2443 ) -> Result<Vec<GoodbyeTarget>, ForwardingError> {
2444 self.cleanup_connection_counted(connection_id)
2445 .map(|cleanup| cleanup.released)
2446 }
2447
2448 pub(crate) fn cleanup_connection_counted(
2453 &self,
2454 connection_id: ConnectionId,
2455 ) -> Result<ConnectionCleanup, ForwardingError> {
2456 let mut inner = self.write_inner()?;
2457 inner.closing_connections.insert(connection_id);
2458 let cleanup = if let Some(endpoint) = inner.endpoint_by_connection.remove(&connection_id) {
2459 remove_module_connection_locked(&mut inner, endpoint)
2460 } else {
2461 ConnectionCleanup {
2462 released: Self::cleanup_client_connection_locked(&mut inner, connection_id),
2463 abandoned_relays: 0,
2464 }
2465 };
2466 inner.closing_connections.remove(&connection_id);
2475 Ok(cleanup)
2476 }
2477
2478 fn cleanup_client_connection_locked(
2479 inner: &mut ForwardingInner,
2480 connection_id: ConnectionId,
2481 ) -> Vec<GoodbyeTarget> {
2482 let routes = inner
2483 .client_to_module
2484 .iter()
2485 .filter(|(key, _)| key.connection_id == connection_id)
2486 .map(|(key, route)| (*key, route.client_epoch))
2487 .collect::<Vec<_>>();
2488 let mut released = Vec::with_capacity(routes.len());
2489 for (client_key, epoch) in routes {
2490 if let RouteRelease::Removed(target) =
2491 release_client_route_locked(inner, client_key, epoch)
2492 {
2493 released.push(target);
2494 }
2495 }
2496
2497 let pending_keys = inner
2498 .pending_relays
2499 .iter()
2500 .filter(|(_, pending)| pending.reservation.client_key.connection_id == connection_id)
2501 .map(|(key, _)| *key)
2502 .collect::<Vec<_>>();
2503 for key in pending_keys {
2504 let Some(pending) = inner.pending_relays.remove(&key) else {
2505 continue;
2506 };
2507 release_reserved_route_locked(
2508 inner,
2509 pending.reservation.client_key,
2510 pending.reservation.module_key,
2511 );
2512 if pending.relay_enqueued {
2513 if let Some(target) = abandoned_route_target(inner, &pending.reservation) {
2514 released.push(target);
2515 }
2516 }
2517 let _ = pending.sender.send(RouteBindRelayOutcome::ModuleGone(
2518 "client connection closed during route.bind relay".to_string(),
2519 ));
2520 }
2521
2522 let orphaned = inner
2523 .reserved_client
2524 .iter()
2525 .filter(|(key, _)| key.connection_id == connection_id)
2526 .map(|(client, module)| (*client, *module))
2527 .collect::<Vec<_>>();
2528 for (client_key, module_key) in orphaned {
2529 release_reserved_route_locked(inner, client_key, module_key);
2530 }
2531 inner.next_client_channel.remove(&connection_id);
2532 inner
2533 .client_slot_epochs
2534 .retain(|key, _| key.connection_id != connection_id);
2535 inner
2536 .last_published_epoch
2537 .retain(|key, _| key.connection_id != connection_id);
2538 inner
2539 .status
2540 .retain(|(key, _), _| key.connection_id != connection_id);
2541
2542 released
2543 }
2544
2545 pub(crate) fn escalate_client_delivery_failure(
2554 &self,
2555 connection_id: ConnectionId,
2556 channel: u16,
2557 expected_epoch: u32,
2558 reason: CloseReason,
2559 undelivered: UndeliveredFrame<'_>,
2560 ) -> Result<bool, ForwardingError> {
2561 let principals = {
2562 let mut inner = self.write_inner()?;
2563 let key = ClientRouteKey {
2564 connection_id,
2565 channel,
2566 };
2567 if inner.last_published_epoch.get(&key).copied() != Some(expected_epoch) {
2568 None
2569 } else {
2570 inner.closing_connections.insert(connection_id);
2571 Some(connection_principals_locked(&inner, connection_id))
2572 }
2573 };
2574 let Some(principals) = principals else {
2575 return Ok(false);
2576 };
2577 let backlog = undelivered.sink.backlog();
2578 let close_reason = reason.to_string();
2579 if self.request_connection_close(connection_id, reason) {
2580 warn!(
2581 connection_id = connection_id.get(),
2582 principals = %principals,
2583 module_id = undelivered.module_id.unwrap_or("unknown"),
2584 client_channel = channel,
2585 queued_bytes = backlog.queued_bytes,
2586 queued_frames = backlog.queued_frames,
2587 oldest_queued_ms = backlog
2588 .oldest_age
2589 .map(|age| age.as_millis() as u64)
2590 .unwrap_or(0),
2591 close_reason = %close_reason,
2592 "closing client connection: its egress queue could not take a frame"
2593 );
2594 }
2595 Ok(true)
2596 }
2597
2598 pub(crate) fn publish_scope_changes(
2612 &self,
2613 changes: &[ScopeTagChange],
2614 ) -> Result<Vec<ScopeDrainedRoute>, ForwardingError> {
2615 let mut inner = self.write_inner()?;
2616 let mut by_scope: HashMap<(&str, &str), &ScopeTagChange> = HashMap::new();
2617 for change in changes {
2618 let key = (change.owner.clone(), change.scope_ref.clone());
2619 match change.after {
2620 Some(tag) => {
2621 inner.scope_tags.insert(key, tag);
2622 }
2623 None => {
2624 inner.scope_tags.remove(&key);
2625 }
2626 }
2627 by_scope.insert((change.owner.as_str(), change.scope_ref.as_str()), change);
2628 }
2629 let mut selected = Vec::new();
2630 for (client_key, route) in &inner.client_to_module {
2631 let Some(scope) = &route.scope else {
2632 continue;
2633 };
2634 let Some(change) = by_scope.get(&(scope.owner.as_str(), scope.scope_ref.as_str()))
2635 else {
2636 continue;
2637 };
2638 if change.before.map(|tag| tag.scope_epoch) != Some(scope.tag.scope_epoch) {
2639 continue;
2640 }
2641 let reason = match &change.drain {
2642 ScopeDrain::Nothing => continue,
2643 ScopeDrain::All(reason) => *reason,
2644 ScopeDrain::Carriers(narrowed) => {
2645 let owner = Principal::Reserved {
2646 module_id: scope.owner.clone(),
2647 };
2648 let hit = route.principal != owner
2649 && narrowed.iter().any(|(principal, allowed)| {
2650 *principal == route.principal
2651 && allowed
2652 .as_ref()
2653 .is_none_or(|targets| !targets.contains(&route.module_id))
2654 });
2655 if !hit {
2656 continue;
2657 }
2658 RouteCloseReason::ScopeCarrierRemoved
2659 }
2660 };
2661 selected.push((*client_key, route.client_epoch, reason));
2662 }
2663 let mut drained = Vec::new();
2664 for (client_key, client_epoch, reason) in selected {
2665 let Some(route) = inner.client_to_module.get(&client_key).cloned() else {
2666 continue;
2667 };
2668 let release = release_client_route_locked(&mut inner, client_key, client_epoch);
2669 self.record_route_release(&release);
2670 if let RouteRelease::Removed(module) = release {
2671 drained.push(ScopeDrainedRoute {
2672 reason,
2673 module_id: route.module_id.clone(),
2674 client: GoodbyeTarget {
2675 connection_id: route.client_connection_id,
2676 sink: route.client_sink.clone(),
2677 negotiated_ver: route.client_negotiated_ver,
2678 channel: route.client_channel,
2679 epoch: route.client_epoch,
2680 kind: GoodbyeTargetKind::Client,
2681 module_id: Some(route.module_id.clone()),
2682 },
2683 module,
2684 });
2685 }
2686 }
2687 Ok(drained)
2688 }
2689
2690 #[cfg(test)]
2692 pub(crate) fn published_scope_tag(&self, owner: &str, scope_ref: &str) -> Option<ScopeTag> {
2693 self.read_inner()
2694 .ok()?
2695 .scope_tags
2696 .get(&(owner.to_string(), scope_ref.to_string()))
2697 .copied()
2698 }
2699
2700 fn record_route_release(&self, release: &RouteRelease) {
2701 match release {
2702 RouteRelease::Removed(_) => self.counters.increment_route_released_epoch_fenced(),
2703 RouteRelease::Stale => self.counters.increment_route_release_stale_skipped(),
2704 RouteRelease::Absent => {}
2705 }
2706 }
2707
2708 fn read_inner(&self) -> Result<RwLockReadGuard<'_, ForwardingInner>, ForwardingError> {
2709 self.inner.read().map_err(|_| ForwardingError::Poisoned)
2710 }
2711
2712 fn write_inner(&self) -> Result<RwLockWriteGuard<'_, ForwardingInner>, ForwardingError> {
2713 self.inner.write().map_err(|_| ForwardingError::Poisoned)
2714 }
2715
2716 fn lock_close_registry(
2717 &self,
2718 ) -> MutexGuard<'_, HashMap<ConnectionId, oneshot::Sender<CloseReason>>> {
2719 self.close_registry
2720 .lock()
2721 .unwrap_or_else(|poisoned| poisoned.into_inner())
2722 }
2723}
2724
2725impl ForwardingInner {
2726 fn allocate_route_slots(
2727 &mut self,
2728 connection_id: ConnectionId,
2729 endpoint: ModuleEndpointId,
2730 ) -> Result<(u16, u32, u16, u32), ForwardingError> {
2731 let client_start = *self.next_client_channel.entry(connection_id).or_insert(1);
2732 let mut client_channel = client_start;
2733 let client_channel = loop {
2734 let key = ClientRouteKey {
2735 connection_id,
2736 channel: client_channel,
2737 };
2738 let eligible = !self.client_to_module.contains_key(&key)
2739 && !self.reserved_client.contains_key(&key)
2740 && self.client_slot_epochs.get(&key).copied().unwrap_or(0) < u32::MAX;
2741 if eligible {
2742 break client_channel;
2743 }
2744 client_channel = next_channel(client_channel);
2745 if client_channel == client_start {
2746 return Err(ForwardingError::ClientRouteChannelExhausted { connection_id });
2747 }
2748 };
2749
2750 let module_start = *self.next_module_channel.entry(endpoint).or_insert(1);
2751 let mut module_channel = module_start;
2752 let module_channel = loop {
2753 let key = ModuleRouteKey {
2754 endpoint,
2755 channel: module_channel,
2756 };
2757 let eligible = !self.module_to_client.contains_key(&key)
2758 && !self.reserved_module.contains_key(&key)
2759 && self.module_slot_epochs.get(&key).copied().unwrap_or(0) < u32::MAX;
2760 if eligible {
2761 break module_channel;
2762 }
2763 module_channel = next_channel(module_channel);
2764 if module_channel == module_start {
2765 return Err(ForwardingError::ModuleRouteChannelExhausted { endpoint });
2766 }
2767 };
2768
2769 let client_key = ClientRouteKey {
2770 connection_id,
2771 channel: client_channel,
2772 };
2773 let module_key = ModuleRouteKey {
2774 endpoint,
2775 channel: module_channel,
2776 };
2777 let client_epoch = self
2778 .client_slot_epochs
2779 .get(&client_key)
2780 .copied()
2781 .unwrap_or(0)
2782 + 1;
2783 let module_epoch = self
2784 .module_slot_epochs
2785 .get(&module_key)
2786 .copied()
2787 .unwrap_or(0)
2788 + 1;
2789 self.client_slot_epochs.insert(client_key, client_epoch);
2790 self.module_slot_epochs.insert(module_key, module_epoch);
2791 self.next_client_channel
2792 .insert(connection_id, next_channel(client_channel));
2793 self.next_module_channel
2794 .insert(endpoint, next_channel(module_channel));
2795 Ok((client_channel, client_epoch, module_channel, module_epoch))
2796 }
2797
2798 fn allocate_control_corr(
2799 &mut self,
2800 endpoint: ModuleEndpointId,
2801 ) -> Result<u64, ForwardingError> {
2802 let candidate = self.next_control_corr.get(&endpoint).copied().unwrap_or(1);
2803 if candidate == 0 {
2804 self.closing_connections.insert(endpoint.connection_id);
2805 return Err(ForwardingError::RelayCorrelationExhausted);
2806 }
2807 self.next_control_corr.insert(
2808 endpoint,
2809 if candidate == u64::MAX {
2810 0
2811 } else {
2812 candidate + 1
2813 },
2814 );
2815 Ok(candidate)
2816 }
2817}
2818
2819fn connection_has_client_routes_locked(
2820 inner: &ForwardingInner,
2821 connection_id: ConnectionId,
2822) -> bool {
2823 inner
2824 .client_to_module
2825 .keys()
2826 .any(|key| key.connection_id == connection_id)
2827 || inner
2828 .reserved_client
2829 .keys()
2830 .any(|key| key.connection_id == connection_id)
2831}
2832
2833fn check_module_connection_role_locked(
2834 inner: &ForwardingInner,
2835 connection_id: ConnectionId,
2836) -> Result<(), ForwardingError> {
2837 if inner.endpoint_by_connection.contains_key(&connection_id)
2838 || connection_has_client_routes_locked(inner, connection_id)
2839 {
2840 return Err(ForwardingError::ConnectionRoleConflict { connection_id });
2841 }
2842 Ok(())
2843}
2844
2845fn next_channel(channel: u16) -> u16 {
2846 let next = channel.wrapping_add(1);
2847 if next == 0 {
2848 1
2849 } else {
2850 next
2851 }
2852}
2853
2854fn endpoint_routes_locked(
2855 inner: &ForwardingInner,
2856 endpoint: ModuleEndpointId,
2857) -> Vec<EndpointRoute> {
2858 let drain_reason = inner.draining_endpoints.get(&endpoint).copied();
2859 let draining = drain_reason.is_some();
2860 let mut routes = inner
2861 .module_to_client
2862 .iter()
2863 .filter(|(module_key, _)| module_key.endpoint == endpoint)
2864 .map(|(_, route)| EndpointRoute {
2865 goodbye_target: GoodbyeTarget {
2866 connection_id: route.client_connection_id,
2867 sink: route.client_sink.clone(),
2868 negotiated_ver: route.client_negotiated_ver,
2869 channel: route.client_channel,
2870 epoch: route.client_epoch,
2871 kind: GoodbyeTargetKind::Client,
2872 module_id: Some(route.module_id.clone()),
2873 },
2874 principal: route.principal.clone(),
2875 bound_at: route.bound_at,
2876 draining,
2877 drain_reason,
2878 })
2879 .collect::<Vec<_>>();
2880 routes.sort_by_key(|route| {
2881 (
2882 route.goodbye_target.connection_id.get(),
2883 route.goodbye_target.channel,
2884 route.goodbye_target.epoch,
2885 )
2886 });
2887 routes
2888}
2889
2890fn scope_refusal_locked(inner: &ForwardingInner, scope: &BoundScope) -> Option<ErrorBody> {
2894 let key = (scope.owner.clone(), scope.scope_ref.clone());
2895 match inner.scope_tags.get(&key) {
2896 Some(current) if *current == scope.tag => None,
2897 Some(current) if current.scope_epoch == scope.tag.scope_epoch => Some(ErrorBody::new(
2898 subc_protocol::error_codes::SCOPE_CHANGED,
2899 format!(
2900 "scope '{}' of {} changed while the route was being bound; re-open it",
2901 scope.scope_ref, scope.owner
2902 ),
2903 )),
2904 _ => Some(ErrorBody::new(
2905 subc_protocol::error_codes::SCOPE_ENDED,
2906 format!(
2907 "scope '{}' of {} at scope_epoch {} ended while the route was being bound",
2908 scope.scope_ref, scope.owner, scope.tag.scope_epoch
2909 ),
2910 )),
2911 }
2912}
2913
2914fn release_reserved_route_locked(
2915 inner: &mut ForwardingInner,
2916 client_key: ClientRouteKey,
2917 module_key: ModuleRouteKey,
2918) {
2919 if inner.reserved_client.get(&client_key).copied() == Some(module_key) {
2920 inner.reserved_client.remove(&client_key);
2921 }
2922 if inner.reserved_module.get(&module_key).copied() == Some(client_key) {
2923 inner.reserved_module.remove(&module_key);
2924 }
2925 inner.status.retain(|(key, _), _| *key != client_key);
2926}
2927
2928fn release_client_route_locked(
2929 inner: &mut ForwardingInner,
2930 client_key: ClientRouteKey,
2931 expected_epoch: u32,
2932) -> RouteRelease {
2933 let Some(route) = inner.client_to_module.get(&client_key) else {
2934 return RouteRelease::Absent;
2935 };
2936 if route.client_epoch != expected_epoch {
2937 return RouteRelease::Stale;
2938 }
2939 let route = inner
2940 .client_to_module
2941 .remove(&client_key)
2942 .expect("route checked under the same forwarding lock");
2943 route.flow.close();
2944 inner.module_to_client.remove(&ModuleRouteKey {
2945 endpoint: route.module_endpoint,
2946 channel: route.module_channel,
2947 });
2948 inner.operator_confirms.route_closed(ModuleRouteKey {
2949 endpoint: route.module_endpoint,
2950 channel: route.module_channel,
2951 });
2952 inner.status.remove(&(client_key, expected_epoch));
2953 RouteRelease::Removed(GoodbyeTarget {
2954 connection_id: route.module_endpoint.connection_id,
2955 sink: route.module_sink.clone(),
2956 negotiated_ver: route.module_negotiated_ver,
2957 channel: route.module_channel,
2958 epoch: route.module_epoch,
2959 kind: GoodbyeTargetKind::Module,
2960 module_id: Some(route.module_id.clone()),
2961 })
2962}
2963
2964fn release_module_route_locked(
2965 inner: &mut ForwardingInner,
2966 module_key: ModuleRouteKey,
2967 expected_epoch: u32,
2968) -> RouteRelease {
2969 let Some(route) = inner.module_to_client.get(&module_key) else {
2970 return RouteRelease::Absent;
2971 };
2972 if route.module_epoch != expected_epoch {
2973 return RouteRelease::Stale;
2974 }
2975 let route = inner
2976 .module_to_client
2977 .remove(&module_key)
2978 .expect("route checked under the same forwarding lock");
2979 inner.operator_confirms.route_closed(module_key);
2980 route.flow.close();
2981 let client_key = ClientRouteKey {
2982 connection_id: route.client_connection_id,
2983 channel: route.client_channel,
2984 };
2985 inner.client_to_module.remove(&client_key);
2986 inner.status.remove(&(client_key, route.client_epoch));
2987 RouteRelease::Removed(GoodbyeTarget {
2988 connection_id: route.client_connection_id,
2989 sink: route.client_sink.clone(),
2990 negotiated_ver: route.client_negotiated_ver,
2991 channel: route.client_channel,
2992 epoch: route.client_epoch,
2993 kind: GoodbyeTargetKind::Client,
2994 module_id: Some(route.module_id.clone()),
2995 })
2996}
2997
2998fn commit_route_locked(
2999 inner: &mut ForwardingInner,
3000 pending: PendingRouteBindRelayEntry,
3001) -> Result<Option<GoodbyeTarget>, ForwardingError> {
3002 let reservation = pending.reservation;
3003 if inner
3004 .closing_connections
3005 .contains(&reservation.client_key.connection_id)
3006 {
3007 return Err(ForwardingError::ConnectionClosing {
3008 connection_id: reservation.client_key.connection_id,
3009 });
3010 }
3011 let module_id = inner
3012 .module_id_by_endpoint
3013 .get(&reservation.module_key.endpoint)
3014 .cloned()
3015 .ok_or(ForwardingError::StaleModuleEndpoint)?;
3016 if inner
3017 .draining_endpoints
3018 .contains_key(&reservation.module_key.endpoint)
3019 {
3020 return Err(ForwardingError::ModuleReloading { module_id });
3021 }
3022 if inner.reserved_client.remove(&reservation.client_key) != Some(reservation.module_key)
3023 || inner.reserved_module.remove(&reservation.module_key) != Some(reservation.client_key)
3024 {
3025 return Err(ForwardingError::UnknownReservation {
3026 client_channel: reservation.client_key.channel,
3027 module_channel: reservation.module_key.channel,
3028 });
3029 }
3030 let module = inner
3031 .modules_by_id
3032 .get(&module_id)
3033 .filter(|module| module.endpoint == reservation.module_key.endpoint)
3034 .cloned()
3035 .ok_or(ForwardingError::StaleModuleEndpoint)?;
3036 let binding = Arc::new(RouteBinding {
3037 client_connection_id: reservation.client_key.connection_id,
3038 client_sink: pending.client_sink,
3039 client_negotiated_ver: pending.client_negotiated_ver,
3040 client_channel: reservation.client_key.channel,
3041 client_epoch: reservation.client_epoch,
3042 module_id,
3043 module_endpoint: reservation.module_key.endpoint,
3044 module_sink: module.sink,
3045 module_negotiated_ver: module.negotiated_ver,
3046 module_channel: reservation.module_key.channel,
3047 module_epoch: reservation.module_epoch,
3048 principal: pending.principal,
3049 project_root: reservation.project_root.clone(),
3050 bound_at: Instant::now(),
3051 flow: Arc::new(ChannelFlow::new(window_for(&module.concurrency))),
3052 scope: pending.scope,
3053 });
3054 inner
3055 .client_to_module
3056 .insert(reservation.client_key, Arc::clone(&binding));
3057 inner
3058 .module_to_client
3059 .insert(reservation.module_key, binding);
3060 let previous_published = inner
3061 .last_published_epoch
3062 .insert(reservation.client_key, reservation.client_epoch);
3063
3064 let client_writer_closed = pending.client_permit.send(pending.route_open_frame);
3069 if client_writer_closed {
3070 let abandoned = pending
3071 .relay_enqueued
3072 .then(|| abandoned_route_target(inner, &reservation))
3073 .flatten();
3074 if let Some(route) = inner.client_to_module.remove(&reservation.client_key) {
3075 route.flow.close();
3076 }
3077 inner.module_to_client.remove(&reservation.module_key);
3078 inner
3079 .status
3080 .remove(&(reservation.client_key, reservation.client_epoch));
3081 match previous_published {
3082 Some(epoch) => {
3083 inner
3084 .last_published_epoch
3085 .insert(reservation.client_key, epoch);
3086 }
3087 None => {
3088 inner.last_published_epoch.remove(&reservation.client_key);
3089 }
3090 }
3091 let _ = pending.sender.send(RouteBindRelayOutcome::ModuleGone(
3092 "client egress closed during route publication".to_string(),
3093 ));
3094 return Ok(abandoned);
3095 }
3096
3097 let _ = pending.sender.send(RouteBindRelayOutcome::Accepted);
3098 Ok(None)
3099}
3100
3101fn module_connection_for_endpoint_locked(
3109 inner: &ForwardingInner,
3110 endpoint: ModuleEndpointId,
3111) -> Option<&ModuleConnection> {
3112 let module_id = inner.module_id_by_endpoint.get(&endpoint)?;
3113 inner
3114 .modules_by_id
3115 .get(module_id)
3116 .filter(|module| module.endpoint == endpoint)
3117 .or_else(|| {
3118 inner
3119 .candidates_by_id
3120 .get(module_id)
3121 .filter(|module| module.endpoint == endpoint)
3122 })
3123 .or_else(|| inner.superseded_endpoints.get(&endpoint))
3124}
3125
3126fn abandoned_route_target(
3127 inner: &ForwardingInner,
3128 reservation: &RouteReservation,
3129) -> Option<GoodbyeTarget> {
3130 let module_id = inner
3131 .module_id_by_endpoint
3132 .get(&reservation.module_key.endpoint)?;
3133 let module = module_connection_for_endpoint_locked(inner, reservation.module_key.endpoint)?;
3134 (module.endpoint == reservation.module_key.endpoint).then(|| GoodbyeTarget {
3135 connection_id: module.endpoint.connection_id,
3136 sink: module.sink.clone(),
3137 negotiated_ver: module.negotiated_ver,
3138 channel: reservation.module_key.channel,
3139 epoch: reservation.module_epoch,
3140 kind: GoodbyeTargetKind::Module,
3141 module_id: Some(module_id.clone()),
3142 })
3143}
3144
3145fn enqueue_hello_ack_locked(
3150 sink: &FrameSink,
3151 connection_id: ConnectionId,
3152 hello_ack: Option<Frame>,
3153) -> Result<(), ForwardingError> {
3154 let Some(hello_ack) = hello_ack else {
3155 return Ok(());
3156 };
3157 sink.try_send(hello_ack)
3158 .map_err(|_| ForwardingError::ModuleEgressUnavailable { connection_id })
3159}
3160
3161fn remove_module_connection_locked(
3162 inner: &mut ForwardingInner,
3163 endpoint: ModuleEndpointId,
3164) -> ConnectionCleanup {
3165 inner.operator_confirms.module_closed(endpoint);
3168 inner.draining_endpoints.remove(&endpoint);
3169 let module_id = inner.module_id_by_endpoint.remove(&endpoint);
3170 if let Some(module_id) = module_id.as_ref() {
3171 if inner
3172 .modules_by_id
3173 .get(module_id)
3174 .is_some_and(|module| module.endpoint == endpoint)
3175 {
3176 inner.modules_by_id.remove(module_id);
3177 }
3178 if inner
3179 .candidates_by_id
3180 .get(module_id)
3181 .is_some_and(|module| module.endpoint == endpoint)
3182 {
3183 inner.candidates_by_id.remove(module_id);
3184 }
3185 }
3186 inner.superseded_endpoints.remove(&endpoint);
3187 inner.endpoint_by_connection.remove(&endpoint.connection_id);
3188 inner.next_module_channel.remove(&endpoint);
3189 inner.next_control_corr.remove(&endpoint);
3190 inner
3191 .health_probe_tombstones
3192 .retain(|(pending_endpoint, _), _| *pending_endpoint != endpoint);
3193 inner
3194 .module_slot_epochs
3195 .retain(|key, _| key.endpoint != endpoint);
3196 let reserved_module_keys: Vec<ModuleRouteKey> = inner
3197 .reserved_module
3198 .keys()
3199 .filter(|module_key| module_key.endpoint == endpoint)
3200 .copied()
3201 .collect();
3202 for module_key in reserved_module_keys {
3203 if let Some(client_key) = inner.reserved_module.get(&module_key).copied() {
3204 release_reserved_route_locked(inner, client_key, module_key);
3205 }
3206 }
3207
3208 let pending_keys: Vec<_> = inner
3209 .pending_relays
3210 .keys()
3211 .filter(|(pending_endpoint, _)| *pending_endpoint == endpoint)
3212 .copied()
3213 .collect();
3214 let pending: Vec<_> = pending_keys
3215 .into_iter()
3216 .filter_map(|key| inner.pending_relays.remove(&key))
3217 .collect();
3218 let abandoned_relays = u32::try_from(pending.len()).unwrap_or(u32::MAX);
3219 for pending in pending {
3220 let module_label = module_id.as_deref().unwrap_or("unknown");
3221 let _ = pending
3222 .sender
3223 .send(RouteBindRelayOutcome::ModuleGone(format!(
3224 "module '{module_label}' connection closed during route.bind relay"
3225 )));
3226 }
3227
3228 let pending_control_keys: Vec<_> = inner
3229 .pending_control_rpcs
3230 .keys()
3231 .filter(|(pending_endpoint, _)| *pending_endpoint == endpoint)
3232 .copied()
3233 .collect();
3234 let pending_control: Vec<_> = pending_control_keys
3235 .into_iter()
3236 .filter_map(|key| inner.pending_control_rpcs.remove(&key))
3237 .collect();
3238 for pending in pending_control {
3239 let module_label = module_id.as_deref().unwrap_or("unknown");
3240 let _ = pending
3241 .sender
3242 .send(ModuleControlRpcOutcome::ModuleGone(format!(
3243 "module '{module_label}' connection closed during module-control RPC"
3244 )));
3245 }
3246
3247 let module_routes = inner
3248 .module_to_client
3249 .iter()
3250 .filter(|(module_key, _)| module_key.endpoint == endpoint)
3251 .map(|(module_key, route)| (*module_key, route.module_epoch))
3252 .collect::<Vec<_>>();
3253 let mut released = Vec::with_capacity(module_routes.len());
3254 for (module_key, epoch) in module_routes {
3255 if let RouteRelease::Removed(target) = release_module_route_locked(inner, module_key, epoch)
3256 {
3257 released.push(target);
3258 }
3259 }
3260 ConnectionCleanup {
3261 released,
3262 abandoned_relays,
3263 }
3264}
3265
3266#[derive(Debug, Clone, Copy)]
3267struct RequestCredit {
3268 subscription: bool,
3269 excluded_from_drain: bool,
3270}
3271
3272#[derive(Debug, Default)]
3273struct CreditLedger {
3274 by_corr: HashMap<u64, Vec<RequestCredit>>,
3275}
3276
3277impl CreditLedger {
3278 fn acquire(&mut self, corr: u64, subscription: bool) {
3279 self.by_corr.entry(corr).or_default().push(RequestCredit {
3280 subscription,
3281 excluded_from_drain: false,
3282 });
3283 }
3284
3285 fn release(&mut self, corr: u64) -> bool {
3286 let Some(credits) = self.by_corr.get_mut(&corr) else {
3287 return false;
3288 };
3289 let released = credits.pop().is_some();
3290 if credits.is_empty() {
3291 self.by_corr.remove(&corr);
3292 }
3293 released
3294 }
3295
3296 fn capture_subscription_exclusions(&mut self) -> u32 {
3297 let mut excluded = 0u32;
3298 for credit in self.by_corr.values_mut().flatten() {
3299 if credit.subscription && !credit.excluded_from_drain {
3300 credit.excluded_from_drain = true;
3301 excluded = excluded.saturating_add(1);
3302 }
3303 }
3304 excluded
3305 }
3306
3307 #[cfg(test)]
3308 fn in_flight(&self) -> usize {
3309 self.by_corr.values().map(Vec::len).sum()
3310 }
3311
3312 fn drain_in_flight(&self) -> usize {
3313 self.by_corr
3314 .values()
3315 .flatten()
3316 .filter(|credit| !credit.excluded_from_drain)
3317 .count()
3318 }
3319
3320 fn drain_held_corrs(&self) -> Vec<u64> {
3323 let mut corrs = self
3324 .by_corr
3325 .iter()
3326 .flat_map(|(corr, credits)| {
3327 credits
3328 .iter()
3329 .filter(|credit| !credit.excluded_from_drain)
3330 .map(move |_| *corr)
3331 })
3332 .collect::<Vec<_>>();
3333 corrs.sort_unstable();
3334 corrs
3335 }
3336}
3337
3338#[derive(Debug, Default)]
3339struct ChannelFlowState {
3340 closed: bool,
3341 credits: CreditLedger,
3342}
3343
3344#[derive(Debug)]
3346pub(crate) struct ChannelFlow {
3347 sem: Semaphore,
3348 window: usize,
3349 state: Mutex<ChannelFlowState>,
3350}
3351
3352impl ChannelFlow {
3353 pub(crate) fn new(window: usize) -> Self {
3354 debug_assert!(window > 0, "flow-control window must be non-zero");
3355 Self {
3356 sem: Semaphore::new(window),
3357 window,
3358 state: Mutex::new(ChannelFlowState::default()),
3359 }
3360 }
3361
3362 #[cfg(test)]
3363 pub(crate) async fn acquire(&self) -> Result<(), ChannelFlowClosed> {
3364 self.acquire_tagged(0, false).await
3365 }
3366
3367 pub(crate) async fn acquire_tagged(
3368 &self,
3369 corr: u64,
3370 subscription: bool,
3371 ) -> Result<(), ChannelFlowClosed> {
3372 let permit = self.sem.acquire().await.map_err(|_| ChannelFlowClosed)?;
3373 let mut state = self
3374 .state
3375 .lock()
3376 .unwrap_or_else(|poisoned| poisoned.into_inner());
3377 if state.closed {
3378 return Err(ChannelFlowClosed);
3379 }
3380 state.credits.acquire(corr, subscription);
3381 permit.forget();
3382 Ok(())
3383 }
3384
3385 #[cfg(test)]
3386 pub(crate) fn release(&self) {
3387 self.release_corr(0);
3388 }
3389
3390 pub(crate) fn release_corr(&self, corr: u64) {
3391 let released = self
3392 .state
3393 .lock()
3394 .unwrap_or_else(|poisoned| poisoned.into_inner())
3395 .credits
3396 .release(corr);
3397 if !released {
3398 warn!(
3402 window = self.window,
3403 available = self.sem.available_permits(),
3404 "flow-control over-release ignored"
3405 );
3406 return;
3407 }
3408 if !self.sem.is_closed() {
3409 self.sem.add_permits(1);
3410 }
3411 }
3412
3413 #[cfg(test)]
3414 pub(crate) fn in_flight(&self) -> usize {
3415 self.state
3416 .lock()
3417 .unwrap_or_else(|poisoned| poisoned.into_inner())
3418 .credits
3419 .in_flight()
3420 }
3421
3422 pub(crate) fn drain_in_flight(&self) -> usize {
3423 self.state
3424 .lock()
3425 .unwrap_or_else(|poisoned| poisoned.into_inner())
3426 .credits
3427 .drain_in_flight()
3428 }
3429
3430 pub(crate) fn drain_held_corrs(&self) -> Vec<u64> {
3431 self.state
3432 .lock()
3433 .unwrap_or_else(|poisoned| poisoned.into_inner())
3434 .credits
3435 .drain_held_corrs()
3436 }
3437
3438 #[cfg(test)]
3439 pub(crate) fn available_permits(&self) -> usize {
3440 self.sem.available_permits()
3441 }
3442
3443 pub(crate) fn begin_drain(&self) -> u32 {
3444 let mut state = self
3445 .state
3446 .lock()
3447 .unwrap_or_else(|poisoned| poisoned.into_inner());
3448 state.closed = true;
3449 self.sem.close();
3450 state.credits.capture_subscription_exclusions()
3451 }
3452
3453 pub(crate) fn close(&self) {
3454 self.state
3455 .lock()
3456 .unwrap_or_else(|poisoned| poisoned.into_inner())
3457 .closed = true;
3458 self.sem.close();
3459 }
3460}
3461
3462#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3463pub(crate) struct ChannelFlowClosed;
3464
3465impl fmt::Display for ChannelFlowClosed {
3466 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3467 write!(f, "flow-control window closed")
3468 }
3469}
3470
3471impl Error for ChannelFlowClosed {}
3472
3473fn window_for(concurrency: &Concurrency) -> usize {
3474 match concurrency {
3475 Concurrency::Serial => 1,
3476 Concurrency::ModuleManaged => DEFAULT_MODULE_MANAGED_WINDOW,
3477 Concurrency::StatelessParallel => STATELESS_PARALLEL_WINDOW,
3478 }
3479}
3480
3481#[derive(Debug, Clone, PartialEq, Eq)]
3482pub enum ForwardingError {
3483 ConnectionRoleConflict {
3486 connection_id: ConnectionId,
3487 },
3488 NoModuleConnection,
3489 ModuleReloading {
3490 module_id: String,
3491 },
3492 StaleModuleEndpoint,
3493 UnknownReservation {
3494 client_channel: u16,
3495 module_channel: u16,
3496 },
3497 ClientRouteChannelExhausted {
3498 connection_id: ConnectionId,
3499 },
3500 ModuleRouteChannelExhausted {
3501 endpoint: ModuleEndpointId,
3502 },
3503 RelayCorrelationExhausted,
3504 ConnectionClosing {
3505 connection_id: ConnectionId,
3506 },
3507 ClientEgressClosed {
3508 connection_id: ConnectionId,
3509 },
3510 RouteOpenBuild(String),
3511 CandidateSlotOccupied {
3513 module_id: String,
3514 },
3515 ModuleEgressUnavailable {
3518 connection_id: ConnectionId,
3519 },
3520 Poisoned,
3521}
3522
3523impl fmt::Display for ForwardingError {
3524 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3525 match self {
3526 Self::ConnectionRoleConflict { connection_id } => write!(f, "connection {} already holds an incompatible role", connection_id.get()),
3527 Self::NoModuleConnection => write!(f, "no module connection is registered"),
3528 Self::ModuleReloading { module_id } => {
3529 write!(f, "module_id '{module_id}' is reloading")
3530 }
3531 Self::StaleModuleEndpoint => write!(f, "module connection generation is stale"),
3532 Self::UnknownReservation {
3533 client_channel,
3534 module_channel,
3535 } => write!(
3536 f,
3537 "route reservation client channel {client_channel} / module channel {module_channel} was not found"
3538 ),
3539 Self::ClientRouteChannelExhausted { connection_id } => write!(
3540 f,
3541 "no client route channels are available for connection {}",
3542 connection_id.get()
3543 ),
3544 Self::ModuleRouteChannelExhausted { endpoint } => write!(
3545 f,
3546 "no module route channels are available for endpoint generation {} on connection {}",
3547 endpoint.generation,
3548 endpoint.connection_id.get()
3549 ),
3550 Self::RelayCorrelationExhausted => {
3551 write!(f, "module control correlation ids are exhausted")
3552 }
3553 Self::ConnectionClosing { connection_id } => write!(
3554 f,
3555 "connection {} is closing and cannot accept route allocation",
3556 connection_id.get()
3557 ),
3558 Self::ClientEgressClosed { connection_id } => write!(
3559 f,
3560 "client connection {} egress is closed",
3561 connection_id.get()
3562 ),
3563 Self::RouteOpenBuild(message) => {
3564 write!(f, "failed to prebuild route.open response: {message}")
3565 }
3566 Self::CandidateSlotOccupied { module_id } => write!(
3567 f,
3568 "module_id '{module_id}' already has a swap candidate registered"
3569 ),
3570 Self::ModuleEgressUnavailable { connection_id } => write!(
3571 f,
3572 "module connection {} egress is unavailable; HELLO_ACK could not be queued",
3573 connection_id.get()
3574 ),
3575 Self::Poisoned => write!(f, "forwarding table lock was poisoned"),
3576 }
3577 }
3578}
3579
3580impl Error for ForwardingError {}
3581
3582#[cfg(test)]
3583mod tests {
3584 use std::time::Duration;
3585
3586 use super::*;
3587 use tokio::sync::mpsc;
3588
3589 #[test]
3590 fn ordinary_long_running_request_is_not_excluded_from_drain() {
3591 let mut ledger = CreditLedger::default();
3592 ledger.acquire(1, false);
3593
3594 assert_eq!(ledger.capture_subscription_exclusions(), 0);
3595 assert_eq!(ledger.drain_in_flight(), 1);
3596 }
3597
3598 #[test]
3599 fn bit_set_subscription_is_excluded_and_counted() {
3600 let mut ledger = CreditLedger::default();
3601 ledger.acquire(1, true);
3602
3603 assert_eq!(ledger.capture_subscription_exclusions(), 1);
3604 assert_eq!(ledger.drain_in_flight(), 0);
3605 }
3606
3607 #[test]
3608 fn subscription_opened_after_drain_snapshot_is_not_excluded() {
3609 let mut ledger = CreditLedger::default();
3610 ledger.acquire(1, true);
3611 assert_eq!(ledger.capture_subscription_exclusions(), 1);
3612
3613 ledger.acquire(2, true);
3614
3615 assert_eq!(ledger.drain_in_flight(), 1);
3616 }
3617
3618 #[test]
3619 fn drain_with_no_subscriptions_reports_zero_excluded() {
3620 let mut ledger = CreditLedger::default();
3621 assert_eq!(ledger.capture_subscription_exclusions(), 0);
3622 }
3623
3624 fn test_hello_ack(corr: u64) -> Frame {
3625 Frame::build(
3626 FrameType::HelloAck,
3627 Flags::new(false, Priority::Passive, false),
3628 0,
3629 0,
3630 corr,
3631 Vec::new(),
3632 )
3633 .unwrap()
3634 }
3635
3636 #[test]
3640 fn acked_registration_that_cannot_queue_its_hello_ack_inserts_nothing() {
3641 let forwarding = ForwardingTable::default();
3642
3643 let (closed_tx, closed_rx) = mpsc::channel(8);
3644 drop(closed_rx);
3645 let closed = ConnectionId::new(1);
3646 assert_eq!(
3647 forwarding.register_module_connection_acked(
3648 closed,
3649 "closed".to_string(),
3650 2,
3651 Concurrency::ModuleManaged,
3652 FrameSink::new(closed_tx),
3653 test_hello_ack(1),
3654 ),
3655 Err(ForwardingError::ModuleEgressUnavailable {
3656 connection_id: closed
3657 })
3658 );
3659
3660 let (full_tx, _full_rx) = mpsc::channel(1);
3661 let full_sink = FrameSink::new(full_tx);
3662 full_sink.try_send(test_hello_ack(99)).unwrap();
3663 let full = ConnectionId::new(2);
3664 assert_eq!(
3665 forwarding.register_module_connection_acked(
3666 full,
3667 "full".to_string(),
3668 2,
3669 Concurrency::ModuleManaged,
3670 full_sink.clone(),
3671 test_hello_ack(2),
3672 ),
3673 Err(ForwardingError::ModuleEgressUnavailable {
3674 connection_id: full
3675 })
3676 );
3677 assert_eq!(
3678 forwarding.register_candidate_module_connection_acked(
3679 full,
3680 "full".to_string(),
3681 2,
3682 Concurrency::ModuleManaged,
3683 full_sink,
3684 test_hello_ack(3),
3685 ),
3686 Err(ForwardingError::ModuleEgressUnavailable {
3687 connection_id: full
3688 })
3689 );
3690
3691 for (connection, module_id) in [(closed, "closed"), (full, "full")] {
3692 assert_eq!(
3693 forwarding
3694 .module_endpoint_for_connection(connection)
3695 .unwrap(),
3696 None
3697 );
3698 let (client_tx, _client_rx) = mpsc::channel(8);
3699 assert_eq!(
3700 forwarding
3701 .begin_route_bind_relay_for_test(
3702 ConnectionId::new(50),
3703 FrameSink::new(client_tx),
3704 1,
3705 module_id,
3706 )
3707 .err(),
3708 Some(ForwardingError::NoModuleConnection)
3709 );
3710 }
3711 assert!(forwarding.read_inner().unwrap().candidates_by_id.is_empty());
3712 }
3713
3714 #[test]
3717 fn acked_registration_queues_the_hello_ack_first() {
3718 let forwarding = ForwardingTable::default();
3719 let (active_tx, mut active_rx) = mpsc::channel(8);
3720 forwarding
3721 .register_module_connection_acked(
3722 ConnectionId::new(1),
3723 "acked".to_string(),
3724 2,
3725 Concurrency::ModuleManaged,
3726 FrameSink::new(active_tx),
3727 test_hello_ack(11),
3728 )
3729 .unwrap();
3730 let (candidate_tx, mut candidate_rx) = mpsc::channel(8);
3731 forwarding
3732 .register_candidate_module_connection_acked(
3733 ConnectionId::new(2),
3734 "acked".to_string(),
3735 2,
3736 Concurrency::ModuleManaged,
3737 FrameSink::new(candidate_tx),
3738 test_hello_ack(12),
3739 )
3740 .unwrap();
3741
3742 let active_first = active_rx.try_recv().unwrap().frame;
3743 assert_eq!(active_first.header.ty, FrameType::HelloAck);
3744 assert_eq!(active_first.header.corr, 11);
3745 let candidate_first = candidate_rx.try_recv().unwrap().frame;
3746 assert_eq!(candidate_first.header.ty, FrameType::HelloAck);
3747 assert_eq!(candidate_first.header.corr, 12);
3748 }
3749
3750 #[test]
3751 fn multi_provider_route_limit_reports_per_client_exhaustion_without_affecting_second_client() {
3752 let forwarding = ForwardingTable::default();
3753 let module_connection = ConnectionId::new(10);
3754 let exhausted_client = ConnectionId::new(20);
3755 let second_client = ConnectionId::new(30);
3756 let (module_tx, _module_rx) = mpsc::channel(1);
3757 let endpoint = forwarding
3758 .register_module_connection(
3759 module_connection,
3760 "route-limit-provider".to_string(),
3761 1,
3762 Concurrency::ModuleManaged,
3763 FrameSink::new(module_tx),
3764 )
3765 .unwrap();
3766
3767 {
3768 let mut inner = forwarding.inner.write().unwrap();
3769 for channel in 1..=u16::MAX {
3770 inner.reserved_client.insert(
3771 ClientRouteKey {
3772 connection_id: exhausted_client,
3773 channel,
3774 },
3775 ModuleRouteKey {
3776 endpoint,
3777 channel: 1,
3778 },
3779 );
3780 }
3781 }
3782
3783 let (exhausted_tx, _exhausted_rx) = mpsc::channel(1);
3784 let err = forwarding
3785 .begin_route_bind_relay_for_test(
3786 exhausted_client,
3787 FrameSink::new(exhausted_tx),
3788 1,
3789 "route-limit-provider",
3790 )
3791 .unwrap_err();
3792 assert!(matches!(
3793 err,
3794 ForwardingError::ClientRouteChannelExhausted { connection_id }
3795 if connection_id == exhausted_client
3796 ));
3797
3798 let (second_tx, _second_rx) = mpsc::channel(1);
3799 let pending = forwarding
3800 .begin_route_bind_relay_for_test(
3801 second_client,
3802 FrameSink::new(second_tx),
3803 2,
3804 "route-limit-provider",
3805 )
3806 .unwrap();
3807 assert_eq!(pending.client_channel, 1);
3808 }
3809
3810 #[test]
3811 fn released_module_channels_are_reused_after_wrap_without_slot_leak() {
3812 let forwarding = ForwardingTable::default();
3813 let module_connection = ConnectionId::new(40);
3814 let client = ConnectionId::new(50);
3815 let (module_tx, _module_rx) = mpsc::channel(1);
3816 forwarding
3817 .register_module_connection(
3818 module_connection,
3819 "slot-reuse-provider".to_string(),
3820 1,
3821 Concurrency::ModuleManaged,
3822 FrameSink::new(module_tx),
3823 )
3824 .unwrap();
3825
3826 let (client_tx, _client_rx) = mpsc::channel(1);
3827 let client_sink = FrameSink::new(client_tx);
3828 let mut wrapped_channel = None;
3829 for index in 0..=usize::from(u16::MAX) {
3830 let pending = forwarding
3831 .begin_route_bind_relay_for_test(
3832 client,
3833 client_sink.clone(),
3834 index as u64 + 1,
3835 "slot-reuse-provider",
3836 )
3837 .unwrap();
3838 if index == usize::from(u16::MAX) {
3839 wrapped_channel = Some(pending.module_channel);
3840 }
3841 forwarding
3842 .abort_pending_relay(
3843 pending.endpoint,
3844 pending.corr,
3845 RouteBindRelayOutcome::ModuleGone("test abort".to_string()),
3846 )
3847 .unwrap();
3848 }
3849
3850 assert_eq!(wrapped_channel, Some(1));
3851 }
3852
3853 #[test]
3854 fn cleanup_connection_prunes_stale_next_client_channel_cursor() {
3855 let forwarding = ForwardingTable::default();
3856 let client = ConnectionId::new(60);
3857 forwarding
3858 .inner
3859 .write()
3860 .unwrap()
3861 .next_client_channel
3862 .insert(client, 41);
3863
3864 let released = forwarding.cleanup_connection(client).unwrap();
3865
3866 assert!(released.is_empty());
3867 assert!(!forwarding
3868 .inner
3869 .read()
3870 .unwrap()
3871 .next_client_channel
3872 .contains_key(&client));
3873 }
3874
3875 #[test]
3876 fn stale_module_cleanup_preserves_fast_reconnect_successor() {
3877 let forwarding = ForwardingTable::default();
3878 let module_id = "fast-reconnect-provider";
3879 let first_connection = ConnectionId::new(70);
3880 let second_connection = ConnectionId::new(80);
3881 let (first_tx, _first_rx) = mpsc::channel(1);
3882 let first_endpoint = forwarding
3883 .register_module_connection(
3884 first_connection,
3885 module_id.to_string(),
3886 1,
3887 Concurrency::ModuleManaged,
3888 FrameSink::new(first_tx),
3889 )
3890 .unwrap();
3891 let (second_tx, _second_rx) = mpsc::channel(1);
3892 let second_endpoint = forwarding
3893 .register_module_connection(
3894 second_connection,
3895 module_id.to_string(),
3896 1,
3897 Concurrency::ModuleManaged,
3898 FrameSink::new(second_tx),
3899 )
3900 .unwrap();
3901 assert_ne!(first_endpoint, second_endpoint);
3902
3903 let released = forwarding.cleanup_connection(first_connection).unwrap();
3904
3905 assert!(released.is_empty());
3906 assert_eq!(
3907 forwarding
3908 .inner
3909 .read()
3910 .unwrap()
3911 .modules_by_id
3912 .get(module_id)
3913 .map(|module| module.endpoint),
3914 Some(second_endpoint)
3915 );
3916 assert!(forwarding.has_live_module_connection(module_id).unwrap());
3917 let control_rpc = forwarding
3918 .begin_module_control_rpc_for(
3919 module_id,
3920 "health.check",
3921 Instant::now() + Duration::from_secs(1),
3922 )
3923 .unwrap();
3924 assert_eq!(control_rpc.endpoint, second_endpoint);
3925 }
3926
3927 fn route_fixture(
3928 module_id: &str,
3929 ) -> (
3930 ForwardingTable,
3931 ConnectionId,
3932 ModuleEndpointId,
3933 ConnectionId,
3934 FrameSink,
3935 mpsc::Receiver<crate::router::OutboundFrame>,
3936 ) {
3937 let forwarding = ForwardingTable::default();
3938 let module_connection = ConnectionId::new(100);
3939 let client_connection = ConnectionId::new(200);
3940 let (module_tx, _module_rx) = mpsc::channel(8);
3941 let endpoint = forwarding
3942 .register_module_connection(
3943 module_connection,
3944 module_id.to_string(),
3945 2,
3946 Concurrency::ModuleManaged,
3947 FrameSink::new(module_tx),
3948 )
3949 .unwrap();
3950 let (client_tx, client_rx) = mpsc::channel(8);
3951 (
3952 forwarding,
3953 module_connection,
3954 endpoint,
3955 client_connection,
3956 FrameSink::new(client_tx),
3957 client_rx,
3958 )
3959 }
3960
3961 #[test]
3962 #[cfg(unix)]
3963 fn daemon_drain_gates_current_and_racing_provider_registrations() {
3964 let (forwarding, _, endpoint, _, sink, _) = route_fixture("provider");
3965 assert_eq!(forwarding.begin_daemon_drain().unwrap(), ["provider"]);
3966 assert!(forwarding.endpoint_is_draining(endpoint).unwrap());
3967 assert!(matches!(
3968 forwarding.register_module_connection(
3969 ConnectionId::new(300),
3970 "late-provider".into(),
3971 2,
3972 Concurrency::ModuleManaged,
3973 sink,
3974 ),
3975 Err(ForwardingError::ConnectionClosing { .. })
3976 ));
3977 }
3978
3979 #[test]
3980 #[cfg(unix)]
3981 fn late_bind_ack_during_daemon_drain_settles_without_leaking_or_closing_module() {
3982 let (forwarding, module_connection, endpoint, client_connection, sink, mut rx) =
3983 route_fixture("provider");
3984 let mut pending =
3985 begin_test_route(&forwarding, client_connection, sink.clone(), 1, "provider");
3986 assert_eq!(forwarding.reserved_route_count().unwrap(), (1, 1));
3987 forwarding.begin_daemon_drain().unwrap();
3988 let completion = forwarding
3989 .complete_pending_relay(
3990 module_connection,
3991 pending.corr,
3992 RouteBindRelayOutcome::Accepted,
3993 )
3994 .expect("draining admission is not a fatal module error");
3995 assert!(completion.settled);
3996 let abandoned = completion
3997 .abandoned
3998 .expect("late accepted bind needs a route GOODBYE");
3999 assert_eq!(abandoned.channel, pending.module_channel);
4000 assert!(matches!(pending.receiver.try_recv().unwrap(),
4001 RouteBindRelayOutcome::Rejected(error) if error.code == "module_reloading"));
4002 assert_eq!(forwarding.reserved_route_count().unwrap(), (0, 0));
4003 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4004 assert!(
4005 rx.try_recv().is_err(),
4006 "no route.open may be published during drain"
4007 );
4008 assert_eq!(
4009 forwarding
4010 .module_endpoint_for_connection(module_connection)
4011 .unwrap(),
4012 Some(endpoint)
4013 );
4014 assert!(
4015 !forwarding
4016 .complete_pending_relay(
4017 module_connection,
4018 pending.corr,
4019 RouteBindRelayOutcome::Accepted,
4020 )
4021 .unwrap()
4022 .settled
4023 );
4024 }
4025
4026 fn test_ping(corr: u64) -> Frame {
4027 Frame::build(
4028 FrameType::Ping,
4029 Flags::new(false, Priority::Passive, false),
4030 0,
4031 0,
4032 corr,
4033 Vec::new(),
4034 )
4035 .unwrap()
4036 }
4037
4038 fn begin_test_route(
4039 forwarding: &ForwardingTable,
4040 client_connection: ConnectionId,
4041 client_sink: FrameSink,
4042 corr: u64,
4043 module_id: &str,
4044 ) -> PendingRouteBindRelay {
4045 forwarding
4046 .begin_route_bind_relay_for_test(client_connection, client_sink, corr, module_id)
4047 .unwrap()
4048 }
4049
4050 #[test]
4051 fn module_registration_and_client_reservation_are_mutually_exclusive_in_both_orders() {
4052 for candidate in [false, true] {
4053 for register_first in [false, true] {
4054 let (forwarding, module_connection, _, client_connection, sink, _rx) =
4055 route_fixture("target");
4056 let register = || {
4057 if candidate {
4058 forwarding.register_candidate_module_connection(
4059 client_connection,
4060 "source".into(),
4061 2,
4062 Concurrency::ModuleManaged,
4063 sink.clone(),
4064 )
4065 } else {
4066 forwarding.register_module_connection(
4067 client_connection,
4068 "source".into(),
4069 2,
4070 Concurrency::ModuleManaged,
4071 sink.clone(),
4072 )
4073 }
4074 };
4075 if register_first {
4076 register().unwrap();
4077 assert!(
4078 forwarding
4079 .begin_route_bind_relay_for_test(
4080 client_connection,
4081 sink.clone(),
4082 1,
4083 "target",
4084 )
4085 .is_err(),
4086 "a deferred route.open cannot reserve after HELLO"
4087 );
4088 assert!(forwarding
4089 .cleanup_connection(client_connection)
4090 .unwrap()
4091 .is_empty());
4092 } else {
4093 let pending =
4094 begin_test_route(&forwarding, client_connection, sink.clone(), 1, "target");
4095 assert!(
4096 register().is_err(),
4097 "HELLO cannot register after a route reservation"
4098 );
4099 forwarding
4100 .complete_pending_relay(
4101 module_connection,
4102 pending.corr,
4103 RouteBindRelayOutcome::Accepted,
4104 )
4105 .unwrap();
4106 assert_eq!(
4107 forwarding
4108 .cleanup_connection(client_connection)
4109 .unwrap()
4110 .len(),
4111 1
4112 );
4113 }
4114 assert_eq!(forwarding.reserved_route_count().unwrap(), (0, 0));
4115 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4116 }
4117 }
4118 }
4119
4120 #[tokio::test]
4126 async fn pending_route_open_completes_behind_queued_data_frames() {
4127 assert_eq!(
4128 crate::server::MAX_PENDING_ROUTE_OPENS_PER_CONNECTION,
4129 8,
4130 "the per-connection pending route.open limit is its own constant"
4131 );
4132 let (forwarding, module_connection, _endpoint, client, _unused_sink, _unused_rx) =
4133 route_fixture("open-behind-data");
4134 let (sink, mut client_rx) = crate::server::connection_egress();
4135 const DATA_FRAMES: usize = 1_000;
4136 let data = |corr: u64| {
4137 Frame::build(
4138 FrameType::StreamData,
4139 Flags::new(false, Priority::Interactive, false),
4140 9,
4141 1,
4142 corr,
4143 vec![b'x'; 200],
4144 )
4145 .unwrap()
4146 };
4147 for corr in 0..DATA_FRAMES as u64 {
4148 sink.try_send(data(corr)).unwrap();
4149 }
4150 let data_bytes = DATA_FRAMES * (subc_protocol::HEADER_LEN + 200);
4151 assert_eq!(sink.backlog().queued_bytes, data_bytes);
4152
4153 let pending = tokio::time::timeout(
4154 Duration::from_secs(5),
4155 forwarding.begin_route_bind_relay_for(
4156 client,
4157 sink.clone(),
4158 subc_protocol::PROTOCOL_VERSION,
4159 4_242,
4160 "open-behind-data",
4161 Principal::Direct,
4162 None,
4163 None,
4164 Instant::now() + Duration::from_secs(60),
4165 ),
4166 )
4167 .await
4168 .expect("reserving the route.open slot must not wait behind data frames")
4169 .unwrap();
4170 forwarding
4171 .complete_pending_relay(
4172 module_connection,
4173 pending.corr,
4174 RouteBindRelayOutcome::Accepted,
4175 )
4176 .unwrap();
4177
4178 let backlog = sink.backlog();
4179 assert_eq!(backlog.queued_frames, DATA_FRAMES + 1);
4180 assert!(
4181 backlog.queued_bytes > data_bytes,
4182 "the route.open response must be counted in queued bytes: {backlog:?}"
4183 );
4184 for corr in 0..DATA_FRAMES as u64 {
4185 assert_eq!(client_rx.recv().await.unwrap().header.corr, corr);
4186 }
4187 let open = client_rx.recv().await.unwrap();
4188 assert_eq!(open.header.corr, 4_242);
4189 assert_eq!(open.header.ty, FrameType::Response);
4190 drop(open);
4191 assert_eq!(sink.backlog().queued_bytes, 0);
4192 assert_eq!(sink.backlog().queued_frames, 0);
4193 }
4194
4195 #[tokio::test]
4200 async fn drain_holdouts_count_held_requests_and_name_the_connection() {
4201 let (forwarding, module_connection, endpoint, client, sink, mut client_rx) =
4202 route_fixture("holdouts");
4203 let mut bound = |corr| {
4204 let route = begin_test_route(&forwarding, client, sink.clone(), corr, "holdouts");
4205 forwarding
4206 .complete_pending_relay(
4207 module_connection,
4208 route.corr,
4209 RouteBindRelayOutcome::Accepted,
4210 )
4211 .unwrap();
4212 client_rx.try_recv().unwrap();
4213 match forwarding
4214 .lookup_data_route(client, route.client_channel, route.client_epoch)
4215 .unwrap()
4216 {
4217 DataRoute::Client(DataRouteState::Bound(binding)) => binding,
4218 other => panic!("expected live route, got {other:?}"),
4219 }
4220 };
4221 let holding = bound(61);
4222 let _idle = bound(62);
4223 holding.flow.acquire_tagged(7, false).await.unwrap();
4224 holding.flow.acquire_tagged(2, false).await.unwrap();
4225 holding.flow.acquire_tagged(3, true).await.unwrap();
4226 forwarding
4227 .begin_module_drain("holdouts", RouteCloseReason::Restart)
4228 .unwrap();
4229
4230 let holdouts = forwarding.endpoint_drain_holdouts(endpoint).unwrap();
4231 assert_eq!(
4232 holdouts,
4233 DrainHoldouts {
4234 requests: 2,
4235 routes: 1,
4236 total_routes: 2,
4237 top_connections: vec![(client.get(), 2)],
4238 held: vec![(holding.module_channel, 2), (holding.module_channel, 7)],
4241 }
4242 );
4243 }
4244
4245 #[test]
4246 fn endpoint_routes_keep_goodbye_targets_and_mark_draining_routes() {
4247 let (forwarding, module_connection, endpoint, client, sink, _client_rx) =
4248 route_fixture("census");
4249 let pending = begin_test_route(&forwarding, client, sink, 1, "census");
4250 forwarding
4251 .complete_pending_relay(
4252 module_connection,
4253 pending.corr,
4254 RouteBindRelayOutcome::Accepted,
4255 )
4256 .unwrap();
4257
4258 let routes = forwarding.endpoint_routes(endpoint).unwrap();
4259 assert_eq!(routes.len(), 1);
4260 assert!(matches!(routes[0].principal, Principal::Direct));
4261 assert_eq!(routes[0].goodbye_target.connection_id, client);
4262 assert_eq!(routes[0].goodbye_target.channel, pending.client_channel);
4263 assert_eq!(routes[0].goodbye_target.epoch, pending.client_epoch);
4264 assert!(!routes[0].draining);
4265
4266 forwarding
4267 .begin_module_drain("census", RouteCloseReason::Restart)
4268 .unwrap();
4269 let draining_routes = forwarding.endpoint_routes(endpoint).unwrap();
4270 assert_eq!(draining_routes.len(), 1);
4271 assert!(draining_routes[0].draining);
4272 }
4273
4274 #[test]
4275 fn aborted_reservation_consumes_both_epochs_and_reuse_advances_them() {
4276 let (forwarding, _, endpoint, client, sink, _client_rx) = route_fixture("epoch-abort");
4277 let first = begin_test_route(&forwarding, client, sink.clone(), 1, "epoch-abort");
4278 assert_eq!((first.client_epoch, first.module_epoch), (1, 1));
4279 forwarding
4280 .abort_pending_relay(
4281 first.endpoint,
4282 first.corr,
4283 RouteBindRelayOutcome::ModuleGone("abort".into()),
4284 )
4285 .unwrap();
4286 forwarding.inject_client_slot_epoch(client, first.client_channel, first.client_epoch);
4287 forwarding.inject_module_slot_epoch(endpoint, first.module_channel, first.module_epoch);
4288
4289 let second = begin_test_route(&forwarding, client, sink, 2, "epoch-abort");
4290 assert_eq!(second.client_channel, first.client_channel);
4291 assert_eq!(second.module_channel, first.module_channel);
4292 assert_eq!((second.client_epoch, second.module_epoch), (2, 2));
4293 }
4294
4295 #[test]
4296 fn stale_release_cannot_remove_reused_successor_and_status_is_epoch_fenced() {
4297 let (forwarding, module_connection, endpoint, client, sink, mut client_rx) =
4298 route_fixture("epoch-release");
4299 let first = begin_test_route(&forwarding, client, sink.clone(), 10, "epoch-release");
4300 forwarding
4301 .complete_pending_relay(
4302 module_connection,
4303 first.corr,
4304 RouteBindRelayOutcome::Accepted,
4305 )
4306 .unwrap();
4307 assert_eq!(client_rx.try_recv().unwrap().header.corr, 10);
4308 assert!(matches!(
4309 forwarding
4310 .release_client_route(client, first.client_channel, first.client_epoch)
4311 .unwrap(),
4312 RouteRelease::Removed(_)
4313 ));
4314 forwarding.inject_client_slot_epoch(client, first.client_channel, first.client_epoch);
4315 forwarding.inject_module_slot_epoch(endpoint, first.module_channel, first.module_epoch);
4316
4317 let second = begin_test_route(&forwarding, client, sink, 11, "epoch-release");
4318 forwarding
4319 .complete_pending_relay(
4320 module_connection,
4321 second.corr,
4322 RouteBindRelayOutcome::Accepted,
4323 )
4324 .unwrap();
4325 assert_eq!(client_rx.try_recv().unwrap().header.corr, 11);
4326 assert!(matches!(
4327 forwarding
4328 .release_client_route(client, second.client_channel, first.client_epoch)
4329 .unwrap(),
4330 RouteRelease::Stale
4331 ));
4332 assert!(!forwarding
4333 .cache_status(
4334 endpoint,
4335 second.module_channel,
4336 first.module_epoch,
4337 "stale".into(),
4338 )
4339 .unwrap());
4340 assert!(forwarding
4341 .cache_status(
4342 endpoint,
4343 second.module_channel,
4344 second.module_epoch,
4345 "current".into(),
4346 )
4347 .unwrap());
4348 match forwarding
4349 .route_poll_snapshot(client, second.client_channel, second.client_epoch)
4350 .unwrap()
4351 {
4352 RoutePollSnapshot::Bound { status, .. } => {
4353 assert_eq!(status.as_deref(), Some("current"));
4354 }
4355 RoutePollSnapshot::Absent => panic!("successor binding was removed"),
4356 }
4357 let counters = forwarding.counters().snapshot();
4358 assert_eq!(counters["route_released_epoch_fenced"], 1);
4359 assert_eq!(counters["route_release_stale_skipped"], 1);
4360 }
4361
4362 #[test]
4363 fn max_epoch_reservation_retires_only_that_slot() {
4364 let (forwarding, _, endpoint, client, sink, _client_rx) = route_fixture("epoch-max");
4365 forwarding.inject_client_slot_epoch(client, 7, u32::MAX - 1);
4366 forwarding.inject_module_slot_epoch(endpoint, 9, u32::MAX - 1);
4367 let final_use = begin_test_route(&forwarding, client, sink.clone(), 20, "epoch-max");
4368 assert_eq!(
4369 (final_use.client_channel, final_use.client_epoch),
4370 (7, u32::MAX)
4371 );
4372 assert_eq!(
4373 (final_use.module_channel, final_use.module_epoch),
4374 (9, u32::MAX)
4375 );
4376 forwarding
4377 .abort_pending_relay(
4378 endpoint,
4379 final_use.corr,
4380 RouteBindRelayOutcome::ModuleGone("abort".into()),
4381 )
4382 .unwrap();
4383 forwarding.inject_client_slot_epoch(client, 7, u32::MAX);
4384 forwarding.inject_module_slot_epoch(endpoint, 9, u32::MAX);
4385 let next = begin_test_route(&forwarding, client, sink, 21, "epoch-max");
4386 assert_ne!(next.client_channel, 7);
4387 assert_ne!(next.module_channel, 9);
4388 assert_eq!((next.client_epoch, next.module_epoch), (1, 1));
4389 }
4390
4391 #[test]
4392 fn bind_and_module_control_share_monotonic_corr_and_deadline_arbitration() {
4393 let (forwarding, module_connection, endpoint, client, sink, _client_rx) =
4394 route_fixture("corr-shared");
4395 let bind = begin_test_route(&forwarding, client, sink, 30, "corr-shared");
4396 assert_eq!(bind.corr, 1);
4397 forwarding
4398 .abort_pending_relay(
4399 endpoint,
4400 bind.corr,
4401 RouteBindRelayOutcome::ModuleGone("abort".into()),
4402 )
4403 .unwrap();
4404 let rpc = forwarding
4405 .begin_module_control_rpc_for(
4406 "corr-shared",
4407 "health.check",
4408 Instant::now() - Duration::from_millis(1),
4409 )
4410 .unwrap();
4411 assert_eq!(rpc.corr, 2);
4412 assert_eq!(
4413 forwarding
4414 .complete_module_control_rpc(
4415 module_connection,
4416 rpc.corr,
4417 Some("health.check"),
4418 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
4419 status: subc_protocol::session::HealthStatus::Ok,
4420 detail: None,
4421 metrics: None,
4422 }),
4423 )
4424 .unwrap(),
4425 ModuleControlRpcCompletion::Settled
4426 );
4427 assert!(matches!(
4428 rpc.receiver.blocking_recv().unwrap(),
4429 ModuleControlRpcOutcome::DeadlineElapsed
4430 ));
4431 }
4432
4433 #[tokio::test(start_paused = true)]
4434 async fn health_probe_tombstone_ttl_removes_an_endpoint_that_stops_probing() {
4435 let (forwarding, _, endpoint, _, _, _) = route_fixture("tombstone-ttl");
4436 let probe_started_at = Instant::now();
4437 let rpc = forwarding
4438 .begin_health_probe_rpc_for(
4439 "tombstone-ttl",
4440 "health.check",
4441 probe_started_at,
4442 probe_started_at + Duration::from_secs(5),
4443 )
4444 .unwrap();
4445 assert!(forwarding
4446 .tombstone_health_probe_rpc(endpoint, rpc.corr)
4447 .unwrap());
4448 assert_eq!(forwarding.health_probe_tombstone_count().unwrap(), 1);
4449
4450 tokio::time::advance(HEALTH_PROBE_TOMBSTONE_TTL).await;
4451 tokio::task::yield_now().await;
4452
4453 assert_eq!(forwarding.health_probe_tombstone_count().unwrap(), 0);
4454 }
4455
4456 #[test]
4457 fn correlation_exhaustion_emits_max_once_then_closes_endpoint() {
4458 let (forwarding, _, endpoint, _, _, _) = route_fixture("corr-max");
4459 let mut close = forwarding.register_connection_close(endpoint.connection_id);
4460 forwarding.inject_control_corr(endpoint, u64::MAX);
4461 let final_rpc = forwarding
4462 .begin_module_control_rpc_for(
4463 "corr-max",
4464 "health.check",
4465 Instant::now() + Duration::from_secs(1),
4466 )
4467 .unwrap();
4468 assert_eq!(final_rpc.corr, u64::MAX);
4469 forwarding
4470 .cancel_module_control_rpc(endpoint, final_rpc.corr)
4471 .unwrap();
4472 assert!(matches!(
4473 forwarding.begin_module_control_rpc_for(
4474 "corr-max",
4475 "health.check",
4476 Instant::now() + Duration::from_secs(1),
4477 ),
4478 Err(ForwardingError::RelayCorrelationExhausted)
4479 ));
4480 assert!(close.try_recv().is_ok());
4481 }
4482
4483 #[test]
4484 fn publication_epoch_controls_delivery_failure_escalation() {
4485 fn setup_successor(
4486 commit_successor: Option<bool>,
4487 ) -> (ForwardingTable, ConnectionId, u16, u32) {
4488 let (forwarding, module_connection, endpoint, client, sink, mut client_rx) =
4489 route_fixture("escalation");
4490 let first = begin_test_route(&forwarding, client, sink.clone(), 40, "escalation");
4491 forwarding
4492 .complete_pending_relay(
4493 module_connection,
4494 first.corr,
4495 RouteBindRelayOutcome::Accepted,
4496 )
4497 .unwrap();
4498 client_rx.try_recv().unwrap();
4499 assert!(matches!(
4500 forwarding
4501 .release_client_route(client, first.client_channel, first.client_epoch)
4502 .unwrap(),
4503 RouteRelease::Removed(_)
4504 ));
4505 if let Some(commit_successor) = commit_successor {
4506 forwarding.inject_client_slot_epoch(
4507 client,
4508 first.client_channel,
4509 first.client_epoch,
4510 );
4511 forwarding.inject_module_slot_epoch(
4512 endpoint,
4513 first.module_channel,
4514 first.module_epoch,
4515 );
4516 let successor = begin_test_route(&forwarding, client, sink, 41, "escalation");
4517 if commit_successor {
4518 forwarding
4519 .complete_pending_relay(
4520 module_connection,
4521 successor.corr,
4522 RouteBindRelayOutcome::Accepted,
4523 )
4524 .unwrap();
4525 client_rx.try_recv().unwrap();
4526 } else {
4527 forwarding
4528 .abort_pending_relay(
4529 endpoint,
4530 successor.corr,
4531 RouteBindRelayOutcome::ModuleGone("abort".into()),
4532 )
4533 .unwrap();
4534 }
4535 }
4536 (forwarding, client, first.client_channel, first.client_epoch)
4537 }
4538
4539 let probe_sink = FrameSink::new(mpsc::channel(1).0);
4540 let (no_successor, client, channel, epoch) = setup_successor(None);
4541 let mut close = no_successor.register_connection_close(client);
4542 assert!(no_successor
4543 .escalate_client_delivery_failure(
4544 client,
4545 channel,
4546 epoch,
4547 CloseReason::new("delivery", "failed"),
4548 UndeliveredFrame {
4549 module_id: None,
4550 sink: &probe_sink,
4551 },
4552 )
4553 .unwrap());
4554 assert!(close.try_recv().is_ok());
4555
4556 let (aborted, client, channel, epoch) = setup_successor(Some(false));
4557 let mut close = aborted.register_connection_close(client);
4558 assert!(aborted
4559 .escalate_client_delivery_failure(
4560 client,
4561 channel,
4562 epoch,
4563 CloseReason::new("delivery", "failed"),
4564 UndeliveredFrame {
4565 module_id: None,
4566 sink: &probe_sink,
4567 },
4568 )
4569 .unwrap());
4570 assert!(close.try_recv().is_ok());
4571
4572 let (published, client, channel, epoch) = setup_successor(Some(true));
4573 let mut close = published.register_connection_close(client);
4574 assert!(!published
4575 .escalate_client_delivery_failure(
4576 client,
4577 channel,
4578 epoch,
4579 CloseReason::new("delivery", "stale failure"),
4580 UndeliveredFrame {
4581 module_id: None,
4582 sink: &probe_sink,
4583 },
4584 )
4585 .unwrap());
4586 assert!(close.try_recv().is_err());
4587 }
4588
4589 #[test]
4590 fn route_concentration_separates_client_count_from_routes_per_client() {
4591 let (forwarding, module_connection, _, client, sink, _client_rx) =
4595 route_fixture("concentration");
4596 assert_eq!(forwarding.client_route_concentration().unwrap(), (0, 0));
4597
4598 for corr in [70_u64, 71] {
4599 let pending =
4600 begin_test_route(&forwarding, client, sink.clone(), corr, "concentration");
4601 forwarding
4602 .complete_pending_relay(
4603 module_connection,
4604 pending.corr,
4605 RouteBindRelayOutcome::Accepted,
4606 )
4607 .unwrap();
4608 }
4609
4610 assert_eq!(forwarding.active_binding_count().unwrap(), 2);
4612 assert_eq!(forwarding.client_route_concentration().unwrap(), (1, 2));
4613 }
4614
4615 #[test]
4616 fn cleanup_and_accepted_resolution_have_one_lock_winner() {
4617 let (forwarding, module_connection, _, client, sink, mut client_rx) =
4618 route_fixture("cleanup-race");
4619 let pending = begin_test_route(&forwarding, client, sink, 45, "cleanup-race");
4620 forwarding
4621 .mark_route_bind_relay_enqueued(pending.endpoint, pending.corr)
4622 .unwrap();
4623 let released = forwarding.cleanup_connection(client).unwrap();
4624 assert_eq!(released.len(), 1);
4625 let completion = forwarding
4626 .complete_pending_relay(
4627 module_connection,
4628 pending.corr,
4629 RouteBindRelayOutcome::Accepted,
4630 )
4631 .unwrap();
4632 assert!(!completion.settled);
4633 assert!(client_rx.try_recv().is_err());
4634 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4635
4636 let (forwarding, module_connection, _, client, sink, mut client_rx) =
4637 route_fixture("accepted-race");
4638 let pending = begin_test_route(&forwarding, client, sink, 46, "accepted-race");
4639 forwarding
4640 .complete_pending_relay(
4641 module_connection,
4642 pending.corr,
4643 RouteBindRelayOutcome::Accepted,
4644 )
4645 .unwrap();
4646 assert_eq!(client_rx.try_recv().unwrap().header.corr, 46);
4647 let released = forwarding.cleanup_connection(client).unwrap();
4648 assert_eq!(released.len(), 1);
4649 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4650 }
4651
4652 #[test]
4653 fn drain_marks_block_reservation_commit_and_live_request_admission_until_phase_two() {
4654 let (forwarding, module_connection, _, client, sink, mut client_rx) =
4655 route_fixture("drain-gap");
4656 let live = begin_test_route(&forwarding, client, sink.clone(), 47, "drain-gap");
4657 forwarding
4658 .complete_pending_relay(
4659 module_connection,
4660 live.corr,
4661 RouteBindRelayOutcome::Accepted,
4662 )
4663 .unwrap();
4664 client_rx.try_recv().unwrap();
4665 let binding = match forwarding
4666 .lookup_data_route(client, live.client_channel, live.client_epoch)
4667 .unwrap()
4668 {
4669 DataRoute::Client(DataRouteState::Bound(binding)) => binding,
4670 other => panic!("expected live route, got {other:?}"),
4671 };
4672
4673 let pending = begin_test_route(&forwarding, client, sink.clone(), 48, "drain-gap");
4674 forwarding
4675 .mark_route_bind_relay_enqueued(pending.endpoint, pending.corr)
4676 .unwrap();
4677 let control_rpc = forwarding
4678 .begin_module_control_rpc_for(
4679 "drain-gap",
4680 "health.check",
4681 Instant::now() + Duration::from_secs(1),
4682 )
4683 .unwrap();
4684 let target = forwarding
4685 .begin_module_drain("drain-gap", RouteCloseReason::Reload)
4686 .unwrap()
4687 .unwrap();
4688 assert!(matches!(
4689 control_rpc.receiver.blocking_recv().unwrap(),
4690 ModuleControlRpcOutcome::ModuleGone(_)
4691 ));
4692 assert_eq!(target.abandoned_bindings.len(), 1);
4693 assert!(binding.flow.sem.is_closed());
4694 assert!(
4695 !forwarding
4696 .complete_pending_relay(
4697 module_connection,
4698 pending.corr,
4699 RouteBindRelayOutcome::Accepted,
4700 )
4701 .unwrap()
4702 .settled
4703 );
4704 assert!(matches!(
4705 forwarding.begin_route_bind_relay_for_test(client, sink, 49, "drain-gap"),
4706 Err(ForwardingError::ModuleReloading { .. })
4707 ));
4708 let released = forwarding
4709 .release_module_endpoint_routes(target.endpoint)
4710 .unwrap();
4711 assert_eq!(released.len(), 1);
4712 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4713 }
4714
4715 #[test]
4723 fn accepted_bind_for_a_closing_client_releases_the_route_instead_of_failing_the_module() {
4724 let (forwarding, module_connection, endpoint, client, sink, mut client_rx) =
4725 route_fixture("closing-client");
4726
4727 let live = begin_test_route(&forwarding, client, sink.clone(), 60, "closing-client");
4730 forwarding
4731 .complete_pending_relay(
4732 module_connection,
4733 live.corr,
4734 RouteBindRelayOutcome::Accepted,
4735 )
4736 .unwrap();
4737 client_rx.try_recv().unwrap();
4738
4739 let pending = begin_test_route(&forwarding, client, sink.clone(), 61, "closing-client");
4741 forwarding
4742 .mark_route_bind_relay_enqueued(pending.endpoint, pending.corr)
4743 .unwrap();
4744
4745 assert!(forwarding
4747 .escalate_client_delivery_failure(
4748 client,
4749 live.client_channel,
4750 live.client_epoch,
4751 CloseReason::new(
4752 "module_to_client_delivery_failed",
4753 "client egress refused a module frame",
4754 ),
4755 UndeliveredFrame {
4756 module_id: None,
4757 sink: &sink,
4758 },
4759 )
4760 .unwrap());
4761 assert!(!sink.is_closed());
4762
4763 let completion = forwarding
4764 .complete_pending_relay(
4765 module_connection,
4766 pending.corr,
4767 RouteBindRelayOutcome::Accepted,
4768 )
4769 .expect("a closing client must not turn a module's ack into an error");
4770
4771 assert!(completion.settled);
4772 let abandoned = completion
4773 .abandoned
4774 .expect("the module must be told to drop the binding it just created");
4775 assert_eq!(abandoned.connection_id, module_connection);
4776 assert_eq!(abandoned.channel, pending.module_channel);
4777 assert_eq!(abandoned.epoch, pending.module_epoch);
4778 assert!(matches!(abandoned.kind, GoodbyeTargetKind::Module));
4779 assert!(matches!(
4780 pending.receiver.blocking_recv().unwrap(),
4781 RouteBindRelayOutcome::ModuleGone(_)
4782 ));
4783 assert!(client_rx.try_recv().is_err());
4786 assert_eq!(forwarding.active_binding_count().unwrap(), 1);
4787
4788 assert!(forwarding
4791 .has_live_module_connection("closing-client")
4792 .unwrap());
4793 let cotenant = ConnectionId::new(201);
4794 let (cotenant_tx, mut cotenant_rx) = mpsc::channel(8);
4795 let cotenant_route = begin_test_route(
4796 &forwarding,
4797 cotenant,
4798 FrameSink::new(cotenant_tx),
4799 62,
4800 "closing-client",
4801 );
4802 assert_eq!(cotenant_route.endpoint, endpoint);
4803 forwarding
4804 .complete_pending_relay(
4805 module_connection,
4806 cotenant_route.corr,
4807 RouteBindRelayOutcome::Accepted,
4808 )
4809 .unwrap();
4810 assert_eq!(cotenant_rx.try_recv().unwrap().header.corr, 62);
4811 assert_eq!(forwarding.active_binding_count().unwrap(), 2);
4812 }
4813
4814 #[test]
4815 fn pending_route_permit_is_released_on_rejection_and_abort() {
4816 let forwarding = ForwardingTable::default();
4817 let module_connection = ConnectionId::new(300);
4818 let client = ConnectionId::new(301);
4819 let (module_tx, _module_rx) = mpsc::channel(1);
4820 let endpoint = forwarding
4821 .register_module_connection(
4822 module_connection,
4823 "permit".into(),
4824 2,
4825 Concurrency::ModuleManaged,
4826 FrameSink::new(module_tx),
4827 )
4828 .unwrap();
4829 let (client_tx, mut client_rx) = mpsc::channel(1);
4830 let sink = FrameSink::new(client_tx);
4831 let rejected = begin_test_route(&forwarding, client, sink.clone(), 50, "permit");
4832 assert!(sink.try_send(test_ping(999)).is_err());
4833 forwarding
4834 .complete_pending_relay(
4835 module_connection,
4836 rejected.corr,
4837 RouteBindRelayOutcome::Rejected(ErrorBody {
4838 code: "no".into(),
4839 message: "rejected".into(),
4840 detail: None,
4841 }),
4842 )
4843 .unwrap();
4844 sink.try_send(test_ping(1000)).unwrap();
4845 assert_eq!(client_rx.try_recv().unwrap().header.corr, 1000);
4846
4847 let aborted = begin_test_route(&forwarding, client, sink.clone(), 51, "permit");
4848 assert!(sink.try_send(test_ping(1001)).is_err());
4849 forwarding
4850 .abort_pending_relay(
4851 endpoint,
4852 aborted.corr,
4853 RouteBindRelayOutcome::ModuleGone("abort".into()),
4854 )
4855 .unwrap();
4856 sink.try_send(test_ping(1002)).unwrap();
4857 assert_eq!(client_rx.try_recv().unwrap().header.corr, 1002);
4858
4859 let receiver_closed = begin_test_route(&forwarding, client, sink, 52, "permit");
4860 forwarding
4861 .mark_route_bind_relay_enqueued(endpoint, receiver_closed.corr)
4862 .unwrap();
4863 drop(client_rx);
4864 let completion = forwarding
4865 .complete_pending_relay(
4866 module_connection,
4867 receiver_closed.corr,
4868 RouteBindRelayOutcome::Accepted,
4869 )
4870 .unwrap();
4871 assert!(completion.abandoned.is_some());
4872 assert_eq!(forwarding.active_binding_count().unwrap(), 0);
4873 }
4874
4875 #[test]
4881 fn cleaned_up_connections_do_not_stay_in_the_closing_set() {
4882 let (forwarding, module_connection, _endpoint, _fixture_client, _sink, _rx) =
4883 route_fixture("closing-set-leak");
4884
4885 const CONNECTIONS: u64 = 32;
4886 for index in 0..CONNECTIONS {
4887 let client = ConnectionId::new(1000 + index);
4888 let (client_tx, _client_rx) = mpsc::channel(8);
4889 let route = begin_test_route(
4890 &forwarding,
4891 client,
4892 FrameSink::new(client_tx),
4893 index + 1,
4894 "closing-set-leak",
4895 );
4896 forwarding
4897 .complete_pending_relay(
4898 module_connection,
4899 route.corr,
4900 RouteBindRelayOutcome::Accepted,
4901 )
4902 .unwrap();
4903 forwarding.cleanup_connection(client).unwrap();
4904 }
4905 forwarding.cleanup_connection(module_connection).unwrap();
4906
4907 assert_eq!(forwarding.closing_connection_count().unwrap(), 0);
4908 }
4909
4910 #[test]
4917 fn closing_connection_is_refused_new_work_until_cleanup_completes() {
4918 let (forwarding, module_connection, _endpoint, client, sink, mut client_rx) =
4919 route_fixture("closing-gate");
4920
4921 let live = begin_test_route(&forwarding, client, sink.clone(), 80, "closing-gate");
4924 forwarding
4925 .complete_pending_relay(
4926 module_connection,
4927 live.corr,
4928 RouteBindRelayOutcome::Accepted,
4929 )
4930 .unwrap();
4931 client_rx.try_recv().unwrap();
4932
4933 assert!(forwarding
4936 .escalate_client_delivery_failure(
4937 client,
4938 live.client_channel,
4939 live.client_epoch,
4940 CloseReason::new(
4941 "module_to_client_delivery_failed",
4942 "client egress refused a module frame",
4943 ),
4944 UndeliveredFrame {
4945 module_id: None,
4946 sink: &sink,
4947 },
4948 )
4949 .unwrap());
4950 assert_eq!(forwarding.closing_connection_count().unwrap(), 1);
4951
4952 assert!(matches!(
4954 forwarding.begin_route_bind_relay_for_test(client, sink, 81, "closing-gate"),
4955 Err(ForwardingError::ConnectionClosing { connection_id })
4956 if connection_id == client
4957 ));
4958 let (late_tx, _late_rx) = mpsc::channel(1);
4960 assert!(matches!(
4961 forwarding.register_module_connection(
4962 client,
4963 "late-module".into(),
4964 2,
4965 Concurrency::ModuleManaged,
4966 FrameSink::new(late_tx),
4967 ),
4968 Err(ForwardingError::ConnectionClosing { connection_id })
4969 if connection_id == client
4970 ));
4971
4972 forwarding.cleanup_connection(client).unwrap();
4976 assert_eq!(forwarding.closing_connection_count().unwrap(), 0);
4977 }
4978}
4979
4980#[cfg(test)]
4983mod swap_slot_tests {
4984 use std::time::Duration;
4985
4986 use super::*;
4987 use tokio::sync::mpsc;
4988
4989 const MODULE_ID: &str = "swapped";
4990
4991 struct SwapFixture {
4992 forwarding: ForwardingTable,
4993 incumbent_connection: ConnectionId,
4994 incumbent: ModuleEndpointId,
4995 candidate_connection: ConnectionId,
4996 candidate: ModuleEndpointId,
4997 _module_rxs: Vec<mpsc::Receiver<crate::router::OutboundFrame>>,
4998 }
4999
5000 fn swap_fixture() -> SwapFixture {
5001 let forwarding = ForwardingTable::default();
5002 let incumbent_connection = ConnectionId::new(100);
5003 let candidate_connection = ConnectionId::new(110);
5004 let (incumbent_tx, incumbent_rx) = mpsc::channel(8);
5005 let incumbent = forwarding
5006 .register_module_connection(
5007 incumbent_connection,
5008 MODULE_ID.to_string(),
5009 2,
5010 Concurrency::ModuleManaged,
5011 FrameSink::new(incumbent_tx),
5012 )
5013 .unwrap();
5014 let (candidate_tx, candidate_rx) = mpsc::channel(8);
5015 let candidate = forwarding
5016 .register_candidate_module_connection(
5017 candidate_connection,
5018 MODULE_ID.to_string(),
5019 2,
5020 Concurrency::ModuleManaged,
5021 FrameSink::new(candidate_tx),
5022 )
5023 .unwrap();
5024 SwapFixture {
5025 forwarding,
5026 incumbent_connection,
5027 incumbent,
5028 candidate_connection,
5029 candidate,
5030 _module_rxs: vec![incumbent_rx, candidate_rx],
5031 }
5032 }
5033
5034 fn client(
5035 raw: u64,
5036 ) -> (
5037 ConnectionId,
5038 FrameSink,
5039 mpsc::Receiver<crate::router::OutboundFrame>,
5040 ) {
5041 let (tx, rx) = mpsc::channel(8);
5042 (ConnectionId::new(raw), FrameSink::new(tx), rx)
5043 }
5044
5045 fn committed_endpoints(forwarding: &ForwardingTable) -> Vec<ModuleEndpointId> {
5046 forwarding
5047 .read_inner()
5048 .unwrap()
5049 .client_to_module
5050 .values()
5051 .map(|route| route.module_endpoint)
5052 .collect()
5053 }
5054
5055 #[test]
5056 fn candidate_is_unroutable_until_cutover_and_by_id_lookups_resolve_the_active_slot() {
5057 let fixture = swap_fixture();
5058 let forwarding = &fixture.forwarding;
5059 assert_ne!(fixture.incumbent, fixture.candidate);
5060
5061 assert!(forwarding.has_live_module_connection(MODULE_ID).unwrap());
5063 assert!(!forwarding.module_is_draining(MODULE_ID).unwrap());
5064 let (client_connection, client_sink, _client_rx) = client(200);
5065 let pending = forwarding
5066 .begin_route_bind_relay_for_test(client_connection, client_sink, 1, MODULE_ID)
5067 .unwrap();
5068 assert_eq!(pending.endpoint, fixture.incumbent);
5069 let rpc = forwarding
5070 .begin_module_control_rpc_for(
5071 MODULE_ID,
5072 "health.check",
5073 Instant::now() + Duration::from_secs(1),
5074 )
5075 .unwrap();
5076 assert_eq!(rpc.endpoint, fixture.incumbent);
5077 let census = forwarding.route_census(Some(MODULE_ID)).unwrap();
5078 assert_eq!(census.len(), 1, "the census lists one endpoint per id");
5079
5080 assert_eq!(
5082 forwarding
5083 .module_endpoint_for_connection(fixture.candidate_connection)
5084 .unwrap(),
5085 Some(fixture.candidate)
5086 );
5087 assert_eq!(
5088 forwarding
5089 .module_id_for_connection(fixture.candidate_connection)
5090 .unwrap()
5091 .as_deref(),
5092 Some(MODULE_ID)
5093 );
5094
5095 let (other_tx, _other_rx) = mpsc::channel(1);
5097 assert_eq!(
5098 forwarding.register_candidate_module_connection(
5099 ConnectionId::new(120),
5100 MODULE_ID.to_string(),
5101 2,
5102 Concurrency::ModuleManaged,
5103 FrameSink::new(other_tx),
5104 ),
5105 Err(ForwardingError::CandidateSlotOccupied {
5106 module_id: MODULE_ID.to_string()
5107 })
5108 );
5109 }
5110
5111 #[test]
5115 fn relay_reserved_before_cutover_never_commits_and_later_relays_land_on_the_candidate() {
5116 let fixture = swap_fixture();
5117 let forwarding = &fixture.forwarding;
5118 let (early_client, early_sink, _early_rx) = client(200);
5119 let mut early = forwarding
5120 .begin_route_bind_relay_for_test(early_client, early_sink, 1, MODULE_ID)
5121 .unwrap();
5122 assert_eq!(early.endpoint, fixture.incumbent);
5123 assert!(forwarding
5124 .mark_route_bind_relay_enqueued(early.endpoint, early.corr)
5125 .unwrap());
5126
5127 let cutover = forwarding.cutover_candidate(MODULE_ID).unwrap().unwrap();
5128 assert_eq!(
5129 cutover,
5130 ForwardingCutover {
5131 promoted: fixture.candidate,
5132 incumbent: Some(fixture.incumbent),
5133 }
5134 );
5135
5136 let (late_client, late_sink, _late_rx) = client(201);
5138 let late = forwarding
5139 .begin_route_bind_relay_for_test(late_client, late_sink, 2, MODULE_ID)
5140 .unwrap();
5141 assert_eq!(
5142 late.endpoint, fixture.candidate,
5143 "a route.open after cutover was reserved on the incumbent"
5144 );
5145
5146 let completion = forwarding
5148 .complete_pending_relay(
5149 fixture.incumbent_connection,
5150 early.corr,
5151 RouteBindRelayOutcome::Accepted,
5152 )
5153 .expect("a superseded endpoint's ack is not an error on its connection");
5154 assert!(completion.settled);
5155 assert!(
5156 !committed_endpoints(forwarding).contains(&fixture.incumbent),
5157 "a relay reserved before cutover committed a route on the incumbent"
5158 );
5159 let goodbye = completion
5160 .abandoned
5161 .expect("the incumbent is told to drop the binding it just created");
5162 assert_eq!(goodbye.connection_id, fixture.incumbent_connection);
5163 assert_eq!(goodbye.channel, early.module_channel);
5164 assert_eq!(goodbye.epoch, early.module_epoch);
5165 assert_eq!(goodbye.kind, GoodbyeTargetKind::Module);
5166 match early.receiver.try_recv() {
5167 Ok(RouteBindRelayOutcome::Rejected(body)) => assert_eq!(body.code, "module_reloading"),
5168 other => panic!("expected a retryable module_reloading answer, got {other:?}"),
5169 }
5170 assert!(matches!(
5171 forwarding
5172 .lookup_data_route(early_client, early.client_channel, early.client_epoch)
5173 .unwrap(),
5174 DataRoute::Client(DataRouteState::Absent)
5175 ));
5176
5177 assert_eq!(forwarding.reserved_route_count().unwrap(), (1, 1));
5180 forwarding
5181 .complete_pending_relay(
5182 fixture.candidate_connection,
5183 late.corr,
5184 RouteBindRelayOutcome::Accepted,
5185 )
5186 .unwrap();
5187 assert_eq!(forwarding.reserved_route_count().unwrap(), (0, 0));
5188 assert_eq!(committed_endpoints(forwarding), vec![fixture.candidate]);
5189 }
5190
5191 #[test]
5192 fn endpoint_drain_after_cutover_drains_the_incumbent_not_the_promoted_candidate() {
5193 let fixture = swap_fixture();
5194 let forwarding = &fixture.forwarding;
5195 let (bound_client, bound_sink, _bound_rx) = client(200);
5197 let bound = forwarding
5198 .begin_route_bind_relay_for_test(bound_client, bound_sink, 1, MODULE_ID)
5199 .unwrap();
5200 forwarding
5201 .complete_pending_relay(
5202 fixture.incumbent_connection,
5203 bound.corr,
5204 RouteBindRelayOutcome::Accepted,
5205 )
5206 .unwrap();
5207 let (pending_client, pending_sink, _pending_rx) = client(201);
5208 let mut in_flight = forwarding
5209 .begin_route_bind_relay_for_test(pending_client, pending_sink, 2, MODULE_ID)
5210 .unwrap();
5211 forwarding
5212 .mark_route_bind_relay_enqueued(in_flight.endpoint, in_flight.corr)
5213 .unwrap();
5214
5215 let incumbent = forwarding
5216 .cutover_candidate(MODULE_ID)
5217 .unwrap()
5218 .unwrap()
5219 .incumbent
5220 .unwrap();
5221 let target = forwarding
5222 .begin_endpoint_drain(incumbent, RouteCloseReason::Restart)
5223 .unwrap()
5224 .expect("the superseded incumbent is still registered");
5225
5226 assert_eq!(target.endpoint, fixture.incumbent);
5227 assert!(forwarding.endpoint_is_draining(fixture.incumbent).unwrap());
5228 assert!(!forwarding.endpoint_is_draining(fixture.candidate).unwrap());
5229 assert!(!forwarding.module_is_draining(MODULE_ID).unwrap());
5230 assert_eq!(target.abandoned_bindings.len(), 1);
5231 assert_eq!(
5232 target.abandoned_bindings[0].channel,
5233 in_flight.module_channel
5234 );
5235 assert!(matches!(
5236 in_flight.receiver.try_recv(),
5237 Ok(RouteBindRelayOutcome::Rejected(body)) if body.code == "module_reloading"
5238 ));
5239 assert_eq!(
5240 forwarding.endpoint_routes(fixture.incumbent).unwrap().len(),
5241 1,
5242 "the incumbent's bound route stays until its drain finishes"
5243 );
5244
5245 let (next_client, next_sink, _next_rx) = client(202);
5246 let next = forwarding
5247 .begin_route_bind_relay_for_test(next_client, next_sink, 3, MODULE_ID)
5248 .expect("the promoted candidate keeps accepting routes");
5249 assert_eq!(next.endpoint, fixture.candidate);
5250 }
5251
5252 #[test]
5258 fn stale_endpoint_ack_without_a_promotion_still_fails_as_before() {
5259 let forwarding = ForwardingTable::default();
5260 let first_connection = ConnectionId::new(70);
5261 let (first_tx, _first_rx) = mpsc::channel(8);
5262 forwarding
5263 .register_module_connection(
5264 first_connection,
5265 MODULE_ID.to_string(),
5266 2,
5267 Concurrency::ModuleManaged,
5268 FrameSink::new(first_tx),
5269 )
5270 .unwrap();
5271 let (client_connection, client_sink, _client_rx) = client(200);
5272 let mut pending = forwarding
5273 .begin_route_bind_relay_for_test(client_connection, client_sink, 1, MODULE_ID)
5274 .unwrap();
5275 let (second_tx, _second_rx) = mpsc::channel(8);
5276 forwarding
5277 .register_module_connection(
5278 ConnectionId::new(80),
5279 MODULE_ID.to_string(),
5280 2,
5281 Concurrency::ModuleManaged,
5282 FrameSink::new(second_tx),
5283 )
5284 .unwrap();
5285
5286 assert_eq!(
5287 forwarding
5288 .complete_pending_relay(
5289 first_connection,
5290 pending.corr,
5291 RouteBindRelayOutcome::Accepted
5292 )
5293 .unwrap_err(),
5294 ForwardingError::StaleModuleEndpoint
5295 );
5296 assert!(committed_endpoints(&forwarding).is_empty());
5297 assert_eq!(forwarding.reserved_route_count().unwrap(), (0, 0));
5298 assert!(matches!(
5299 pending.receiver.try_recv(),
5300 Err(oneshot::error::TryRecvError::Closed)
5301 ));
5302 }
5303
5304 #[test]
5305 fn cleanup_releases_candidate_and_superseded_slots_without_touching_the_active_one() {
5306 let fixture = swap_fixture();
5308 let forwarding = &fixture.forwarding;
5309 assert!(forwarding
5310 .cleanup_connection(fixture.candidate_connection)
5311 .unwrap()
5312 .is_empty());
5313 assert_eq!(forwarding.cutover_candidate(MODULE_ID).unwrap(), None);
5314 let (client_connection, client_sink, _client_rx) = client(200);
5315 assert_eq!(
5316 forwarding
5317 .begin_route_bind_relay_for_test(client_connection, client_sink, 1, MODULE_ID)
5318 .unwrap()
5319 .endpoint,
5320 fixture.incumbent
5321 );
5322
5323 let fixture = swap_fixture();
5326 let forwarding = &fixture.forwarding;
5327 let (bound_client, bound_sink, _bound_rx) = client(200);
5328 let bound = forwarding
5329 .begin_route_bind_relay_for_test(bound_client, bound_sink, 1, MODULE_ID)
5330 .unwrap();
5331 forwarding
5332 .complete_pending_relay(
5333 fixture.incumbent_connection,
5334 bound.corr,
5335 RouteBindRelayOutcome::Accepted,
5336 )
5337 .unwrap();
5338 forwarding.cutover_candidate(MODULE_ID).unwrap().unwrap();
5339 let released = forwarding
5340 .cleanup_connection(fixture.incumbent_connection)
5341 .unwrap();
5342 assert_eq!(released.len(), 1);
5343 assert_eq!(released[0].connection_id, bound_client);
5344 assert!(forwarding
5345 .read_inner()
5346 .unwrap()
5347 .superseded_endpoints
5348 .is_empty());
5349 assert!(forwarding.has_live_module_connection(MODULE_ID).unwrap());
5350 let (next_client, next_sink, _next_rx) = client(201);
5351 assert_eq!(
5352 forwarding
5353 .begin_route_bind_relay_for_test(next_client, next_sink, 2, MODULE_ID)
5354 .unwrap()
5355 .endpoint,
5356 fixture.candidate
5357 );
5358 }
5359}