1use std::{
2 collections::HashMap,
3 error::Error,
4 fmt,
5 sync::{
6 atomic::{AtomicU64, AtomicUsize, Ordering},
7 Arc, Mutex,
8 },
9 time::{Duration, Instant},
10};
11
12use subc_protocol::{error_codes, ErrorBody, Flags, FrameType, Priority};
13use tokio::sync::{mpsc, Notify};
14use tracing::debug;
15
16use crate::{
17 control::ControlHandler,
18 forwarding::{
19 CloseReason, ConnectionCloseReceiver, DataRoute, DataRouteState, ForwardingError,
20 ForwardingTable, RouteBinding, RouteRelease, UndeliveredFrame,
21 },
22 registry::ConnectionId,
23 DaemonCounters, Frame, FrameBuildError,
24};
25
26#[derive(Debug)]
35pub struct OutboundFrame {
36 pub frame: Frame,
37 pub enqueued_at: std::time::Instant,
38 pub(crate) flushed: Option<tokio::sync::oneshot::Sender<()>>,
39 pub(crate) charge: Option<EgressCharge>,
44}
45
46impl OutboundFrame {
47 fn charged(frame: Frame, charge: EgressCharge) -> Self {
48 Self {
49 frame,
50 enqueued_at: charge.enqueued_at,
51 flushed: None,
52 charge: Some(charge),
53 }
54 }
55}
56
57fn queued_frame_bytes(frame: &Frame) -> usize {
60 subc_protocol::HEADER_LEN + frame.body.len()
61}
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq)]
66pub(crate) struct EgressBacklog {
67 pub queued_bytes: usize,
68 pub queued_frames: usize,
69 pub oldest_age: Option<Duration>,
76}
77
78#[derive(Debug)]
83struct EgressAccounting {
84 byte_budget: usize,
85 queued_bytes: AtomicUsize,
86 queued_frames: AtomicUsize,
87 time_base: Instant,
89 oldest_enqueued_nanos: AtomicU64,
92 waiters: AtomicUsize,
95 freed: Notify,
98}
99
100impl EgressAccounting {
101 fn new(byte_budget: usize) -> Self {
102 Self {
103 byte_budget,
104 queued_bytes: AtomicUsize::new(0),
105 queued_frames: AtomicUsize::new(0),
106 time_base: Instant::now(),
107 oldest_enqueued_nanos: AtomicU64::new(0),
108 waiters: AtomicUsize::new(0),
109 freed: Notify::new(),
110 }
111 }
112
113 fn stamp(&self, at: Instant) -> u64 {
114 (at.saturating_duration_since(self.time_base).as_nanos() as u64).saturating_add(1)
115 }
116
117 fn try_charge(self: &Arc<Self>, bytes: usize) -> Option<EgressCharge> {
122 let mut current = self.queued_bytes.load(Ordering::SeqCst);
125 loop {
126 if current != 0 && current.saturating_add(bytes) > self.byte_budget {
127 return None;
128 }
129 match self.queued_bytes.compare_exchange_weak(
130 current,
131 current + bytes,
132 Ordering::SeqCst,
133 Ordering::SeqCst,
134 ) {
135 Ok(_) => return Some(self.record(bytes)),
136 Err(actual) => current = actual,
137 }
138 }
139 }
140
141 fn charge_unconditionally(self: &Arc<Self>, bytes: usize) -> EgressCharge {
144 self.queued_bytes.fetch_add(bytes, Ordering::SeqCst);
145 self.record(bytes)
146 }
147
148 fn record(self: &Arc<Self>, bytes: usize) -> EgressCharge {
149 let enqueued_at = Instant::now();
150 let stamp = self.stamp(enqueued_at);
151 if self.queued_frames.fetch_add(1, Ordering::AcqRel) == 0 {
152 self.oldest_enqueued_nanos.store(stamp, Ordering::Release);
155 }
156 EgressCharge {
157 accounting: Arc::clone(self),
158 bytes,
159 stamp,
160 enqueued_at,
161 }
162 }
163
164 fn taken_by_writer(&self, stamp: u64) {
166 self.oldest_enqueued_nanos.store(stamp, Ordering::Release);
167 }
168
169 fn release(&self, bytes: usize, stamp: u64) {
170 self.queued_bytes.fetch_sub(bytes, Ordering::SeqCst);
171 if self.queued_frames.fetch_sub(1, Ordering::AcqRel) == 1 {
172 let _ = self.oldest_enqueued_nanos.compare_exchange(
175 stamp,
176 0,
177 Ordering::AcqRel,
178 Ordering::Acquire,
179 );
180 }
181 if self.waiters.load(Ordering::SeqCst) != 0 {
182 self.freed.notify_waiters();
183 }
184 }
185
186 fn backlog(&self) -> EgressBacklog {
187 let queued_frames = self.queued_frames.load(Ordering::Acquire);
188 let stamp = self.oldest_enqueued_nanos.load(Ordering::Acquire);
189 let oldest_age = (queued_frames != 0 && stamp != 0).then(|| {
190 let enqueued = self.time_base + Duration::from_nanos(stamp - 1);
191 enqueued.elapsed()
192 });
193 EgressBacklog {
194 queued_bytes: self.queued_bytes.load(Ordering::Acquire),
195 queued_frames,
196 oldest_age,
197 }
198 }
199}
200
201#[derive(Debug)]
204pub(crate) struct EgressCharge {
205 accounting: Arc<EgressAccounting>,
206 bytes: usize,
207 stamp: u64,
208 enqueued_at: Instant,
209}
210
211impl EgressCharge {
212 pub(crate) fn taken_by_writer(&self) {
215 self.accounting.taken_by_writer(self.stamp);
216 }
217}
218
219impl Drop for EgressCharge {
220 fn drop(&mut self) {
221 self.accounting.release(self.bytes, self.stamp);
222 }
223}
224
225#[derive(Debug)]
231pub(crate) struct EgressPermit {
232 permit: mpsc::OwnedPermit<OutboundFrame>,
233 accounting: Arc<EgressAccounting>,
234}
235
236impl EgressPermit {
237 pub(crate) fn send(self, frame: Frame) -> bool {
240 let charge = self
241 .accounting
242 .charge_unconditionally(queued_frame_bytes(&frame));
243 let sender = self.permit.send(OutboundFrame::charged(frame, charge));
244 sender.is_closed()
245 }
246}
247
248impl std::ops::Deref for OutboundFrame {
251 type Target = Frame;
252
253 fn deref(&self) -> &Frame {
254 &self.frame
255 }
256}
257
258#[cfg(test)]
261pub(crate) mod test_log {
262 use std::{
263 io::Write,
264 sync::{Arc, Mutex},
265 };
266
267 #[derive(Clone)]
268 struct TestLogWriter(Arc<Mutex<Vec<u8>>>);
269
270 impl Write for TestLogWriter {
271 fn write(&mut self, buffer: &[u8]) -> std::io::Result<usize> {
272 self.0
273 .lock()
274 .expect("test log capture is not poisoned")
275 .extend(buffer);
276 Ok(buffer.len())
277 }
278
279 fn flush(&mut self) -> std::io::Result<()> {
280 Ok(())
281 }
282 }
283
284 pub(crate) fn log_capture(
285 level: tracing::Level,
286 ) -> (Arc<Mutex<Vec<u8>>>, tracing::dispatcher::DefaultGuard) {
287 let output = Arc::new(Mutex::new(Vec::new()));
288 let writer = Arc::clone(&output);
289 let subscriber = tracing_subscriber::fmt()
290 .with_max_level(level)
291 .with_ansi(false)
292 .without_time()
293 .with_target(false)
294 .with_writer(move || TestLogWriter(Arc::clone(&writer)))
295 .finish();
296 let guard = tracing::subscriber::set_default(subscriber);
297 (output, guard)
298 }
299
300 pub(crate) fn captured_logs(output: &Arc<Mutex<Vec<u8>>>) -> String {
301 String::from_utf8(
302 output
303 .lock()
304 .expect("test log capture is not poisoned")
305 .clone(),
306 )
307 .expect("tracing output is UTF-8")
308 }
309}
310
311#[derive(Debug, Clone)]
319pub struct FrameSink {
320 tx: mpsc::Sender<OutboundFrame>,
321 accounting: Arc<EgressAccounting>,
322}
323
324impl FrameSink {
325 pub fn new(tx: mpsc::Sender<OutboundFrame>) -> Self {
329 Self::with_byte_budget(tx, crate::server::CONNECTION_EGRESS_BYTE_BUDGET)
330 }
331
332 pub(crate) fn with_byte_budget(tx: mpsc::Sender<OutboundFrame>, byte_budget: usize) -> Self {
333 Self {
334 tx,
335 accounting: Arc::new(EgressAccounting::new(byte_budget)),
336 }
337 }
338
339 async fn charge_waiting(&self, bytes: usize) -> Option<EgressCharge> {
346 if let Some(charge) = self.accounting.try_charge(bytes) {
347 return Some(charge);
348 }
349 struct Waiting<'a>(&'a AtomicUsize);
352 impl Drop for Waiting<'_> {
353 fn drop(&mut self) {
354 self.0.fetch_sub(1, Ordering::SeqCst);
355 }
356 }
357 self.accounting.waiters.fetch_add(1, Ordering::SeqCst);
358 let _waiting = Waiting(&self.accounting.waiters);
359 loop {
360 let freed = self.accounting.freed.notified();
361 tokio::pin!(freed);
362 freed.as_mut().enable();
365 if let Some(charge) = self.accounting.try_charge(bytes) {
366 return Some(charge);
367 }
368 tokio::select! {
369 _ = &mut freed => {}
370 _ = self.tx.closed() => return None,
371 }
372 }
373 }
374
375 pub async fn send(&self, frame: Frame) -> Result<(), RouterError> {
376 let channel = frame.header.channel;
377 let epoch = frame.header.epoch;
378 let corr = frame.header.corr;
379 let closed =
380 || RouterError::backend_with_epoch(channel, epoch, corr, "connection writer closed");
381 let charge = self
382 .charge_waiting(queued_frame_bytes(&frame))
383 .await
384 .ok_or_else(closed)?;
385 self.tx
386 .send(OutboundFrame::charged(frame, charge))
387 .await
388 .map_err(|_| closed())
389 }
390
391 #[cfg(unix)]
394 pub(crate) async fn send_flushed(&self, frame: Frame) -> Result<(), RouterError> {
395 let (tx, rx) = tokio::sync::oneshot::channel();
396 let charge = self
397 .charge_waiting(queued_frame_bytes(&frame))
398 .await
399 .ok_or_else(|| RouterError::backend(0, 0, "connection writer closed"))?;
400 let mut outbound = OutboundFrame::charged(frame, charge);
401 outbound.flushed = Some(tx);
402 self.tx
403 .send(outbound)
404 .await
405 .map_err(|_| RouterError::backend(0, 0, "connection writer closed"))?;
406 rx.await
407 .map_err(|_| RouterError::backend(0, 0, "connection flush failed"))
408 }
409
410 pub(crate) async fn reserve_owned(&self) -> Result<EgressPermit, RouterError> {
411 let permit = self
412 .tx
413 .clone()
414 .reserve_owned()
415 .await
416 .map_err(|_| RouterError::backend(0, 0, "connection writer closed"))?;
417 Ok(EgressPermit {
418 permit,
419 accounting: Arc::clone(&self.accounting),
420 })
421 }
422
423 #[cfg(test)]
424 pub(crate) fn try_reserve_owned(&self) -> Result<EgressPermit, RouterError> {
425 let permit = self
426 .tx
427 .clone()
428 .try_reserve_owned()
429 .map_err(|err| RouterError::backend(0, 0, err.to_string()))?;
430 Ok(EgressPermit {
431 permit,
432 accounting: Arc::clone(&self.accounting),
433 })
434 }
435
436 pub(crate) fn is_closed(&self) -> bool {
437 self.tx.is_closed()
438 }
439
440 pub(crate) fn backlog(&self) -> EgressBacklog {
442 self.accounting.backlog()
443 }
444
445 pub(crate) fn try_send(&self, frame: Frame) -> Result<(), RouterError> {
449 let channel = frame.header.channel;
450 let epoch = frame.header.epoch;
451 let corr = frame.header.corr;
452 let unavailable = |why: String| {
453 RouterError::backend_with_epoch(
454 channel,
455 epoch,
456 corr,
457 format!("connection writer unavailable: {why}"),
458 )
459 };
460 let bytes = queued_frame_bytes(&frame);
461 let Some(charge) = self.accounting.try_charge(bytes) else {
462 if self.tx.is_closed() {
463 return Err(unavailable("channel closed".to_string()));
464 }
465 return Err(unavailable(format!(
466 "egress byte budget exhausted ({} queued bytes, frame of {bytes} bytes, budget {})",
467 self.accounting.queued_bytes.load(Ordering::Acquire),
468 self.accounting.byte_budget
469 )));
470 };
471 self.tx
473 .try_send(OutboundFrame::charged(frame, charge))
474 .map_err(|err| unavailable(err.to_string()))
475 }
476}
477
478const ORPHAN_ROUTE_GOODBYE_INTERVAL: Duration = Duration::from_secs(5);
490
491const ORPHAN_ROUTE_GOODBYE_PRUNE_AT: usize = 256;
498
499#[derive(Debug, Default)]
503struct OrphanGoodbyeLimiter {
504 last_sent: Mutex<HashMap<ConnectionId, HashMap<u16, tokio::time::Instant>>>,
505}
506
507impl OrphanGoodbyeLimiter {
508 fn claim(&self, connection_id: ConnectionId, channel: u16) -> bool {
511 let now = tokio::time::Instant::now();
512 let mut last_sent = self
513 .last_sent
514 .lock()
515 .unwrap_or_else(std::sync::PoisonError::into_inner);
516 let channels = last_sent.entry(connection_id).or_default();
517 if channels.get(&channel).is_some_and(|sent| {
518 now.saturating_duration_since(*sent) < ORPHAN_ROUTE_GOODBYE_INTERVAL
519 }) {
520 return false;
521 }
522 if channels.len() >= ORPHAN_ROUTE_GOODBYE_PRUNE_AT {
523 channels.retain(|_, sent| {
524 now.saturating_duration_since(*sent) < ORPHAN_ROUTE_GOODBYE_INTERVAL
525 });
526 }
527 channels.insert(channel, now);
528 true
529 }
530
531 fn forget_connection(&self, connection_id: ConnectionId) {
532 self.last_sent
533 .lock()
534 .unwrap_or_else(std::sync::PoisonError::into_inner)
535 .remove(&connection_id);
536 }
537}
538
539#[derive(Debug, Clone)]
541pub struct RouteCtx {
542 pub connection_id: ConnectionId,
543 pub egress: FrameSink,
544}
545
546#[derive(Debug, Clone)]
552pub enum Backend {
553 Echo(EchoBackend),
554 Forward(ForwardBackend),
555}
556
557impl From<EchoBackend> for Backend {
558 fn from(backend: EchoBackend) -> Self {
559 Self::Echo(backend)
560 }
561}
562
563impl From<ForwardBackend> for Backend {
564 fn from(backend: ForwardBackend) -> Self {
565 Self::Forward(backend)
566 }
567}
568
569impl Backend {
570 pub async fn handle(&self, ctx: RouteCtx, frame: Frame) -> Result<(), RouterError> {
571 match self {
572 Self::Echo(backend) => backend.handle(ctx, frame).await,
573 Self::Forward(backend) => backend.handle(ctx, frame).await,
574 }
575 }
576}
577
578pub struct Router {
599 backends: HashMap<u16, Backend>,
600 control: Arc<ControlHandler>,
601 forwarding: Arc<ForwardingTable>,
602 forward_backend: ForwardBackend,
603 counters: DaemonCounters,
604 next_connection_id: AtomicU64,
605 orphan_goodbyes: Arc<OrphanGoodbyeLimiter>,
606}
607
608impl Router {
609 pub fn with_control_handler(control: Arc<ControlHandler>) -> Self {
610 control.install_swap_promotion_observer();
613 let forwarding = control.forwarding();
614 let counters = control.counters();
615 Self {
616 backends: HashMap::new(),
617 control,
618 forwarding: Arc::clone(&forwarding),
619 forward_backend: ForwardBackend::new(forwarding),
620 counters,
621 next_connection_id: AtomicU64::new(1),
623 orphan_goodbyes: Arc::default(),
624 }
625 }
626
627 pub fn with_default_self_handler() -> Self {
628 Self::with_control_handler(Arc::new(ControlHandler::default()))
629 }
630
631 pub fn forwarding(&self) -> Arc<ForwardingTable> {
632 Arc::clone(&self.forwarding)
633 }
634
635 pub fn register_backend(
636 &mut self,
637 channel: u16,
638 backend: impl Into<Backend>,
639 ) -> Result<(), RouterError> {
640 self.register_backend_arc(channel, Arc::new(backend.into()))
641 }
642
643 pub(crate) fn register_backend_arc(
644 &mut self,
645 channel: u16,
646 backend: Arc<Backend>,
647 ) -> Result<(), RouterError> {
648 if channel == 0 {
649 return Err(RouterError::ReservedChannelZero);
650 }
651 if self.backends.contains_key(&channel) {
652 return Err(RouterError::DuplicateChannel { channel });
653 }
654 self.backends.insert(channel, backend.as_ref().clone());
655 Ok(())
656 }
657
658 fn record_module_frame_drop(&self, connection_id: ConnectionId) -> Result<(), RouterError> {
661 let module_id = self
662 .forwarding
663 .module_id_for_connection(connection_id)
664 .map_err(RouterError::Forwarding)?;
665 self.counters
666 .increment_module_frames_dropped_no_route(module_id.as_deref());
667 Ok(())
668 }
669
670 fn handle_orphan_module_frame(&self, ctx: &RouteCtx, frame: &Frame) -> Result<(), RouterError> {
682 let channel = frame.header.channel;
683 let epoch = frame.header.epoch;
684 let module_id = self
685 .forwarding
686 .module_id_for_connection(ctx.connection_id)
687 .map_err(RouterError::Forwarding)?;
688 self.counters
689 .increment_module_frames_dropped_no_route(module_id.as_deref());
690 if self
691 .forwarding
692 .module_route_epoch_was_allocated(ctx.connection_id, channel, epoch)
693 .map_err(RouterError::Forwarding)?
694 {
695 self.counters
696 .increment_module_frames_dropped_released_route(module_id.as_deref());
697 }
698 if frame.header.ty == FrameType::Goodbye
699 || !self.orphan_goodbyes.claim(ctx.connection_id, channel)
700 {
701 return Ok(());
702 }
703 let goodbye = Frame::build_with_version(
704 frame.header.ver,
705 FrameType::Goodbye,
706 Flags::new(false, Priority::Passive, false),
707 channel,
708 epoch,
709 0,
710 Vec::new(),
711 )
712 .map_err(RouterError::FrameBuild)?;
713 match ctx.egress.try_send(goodbye) {
714 Ok(()) => {
715 self.counters.increment_module_orphan_route_goodbyes_sent();
716 debug!(
717 connection_id = ctx.connection_id.get(),
718 module_id = module_id.as_deref().unwrap_or("unknown"),
719 channel,
720 epoch,
721 "answered module frame on a route the daemon does not hold with a route GOODBYE"
722 );
723 }
724 Err(err) => debug!(
725 connection_id = ctx.connection_id.get(),
726 module_id = module_id.as_deref().unwrap_or("unknown"),
727 channel,
728 epoch,
729 error = %err,
730 "could not enqueue route GOODBYE for module frame on a route the daemon does not hold; the next such frame retries"
731 ),
732 }
733 Ok(())
734 }
735
736 pub fn begin_connection(&self) -> RouterConnection {
737 let raw = self.next_connection_id.fetch_add(1, Ordering::Relaxed);
738 let id = ConnectionId::new(raw);
739 let close_receiver = self.forwarding.register_connection_close(id);
740 RouterConnection {
741 id,
742 control_handler: Arc::clone(&self.control),
743 forwarding: Arc::clone(&self.forwarding),
744 close_receiver: Some(close_receiver),
745 orphan_goodbyes: Arc::clone(&self.orphan_goodbyes),
746 }
747 }
748
749 pub(crate) fn route_open_target(&self, frame: &Frame) -> Option<String> {
750 self.control.route_open_target(frame)
751 }
752
753 pub(crate) fn route_open_capacity_refusal(
754 &self,
755 ctx: &RouteCtx,
756 frame: &Frame,
757 target_module_id: &str,
758 in_flight: usize,
759 limit: usize,
760 ) -> Result<Frame, RouterError> {
761 self.control
762 .route_open_capacity_refusal(ctx, frame, target_module_id, in_flight, limit)
763 }
764
765 pub async fn route_for_connection(
766 &self,
767 ctx: &RouteCtx,
768 frame: Frame,
769 ) -> Result<(), RouterError> {
770 self.route_for_connection_started(ctx, frame, None).await
771 }
772
773 pub(crate) async fn route_for_connection_started(
774 &self,
775 ctx: &RouteCtx,
776 frame: Frame,
777 dispatch_started_at: Option<Instant>,
778 ) -> Result<(), RouterError> {
779 let channel = frame.header.channel;
780 let epoch = frame.header.epoch;
781 let corr = frame.header.corr;
782 if channel == 0 {
783 debug!(
784 connection_id = ctx.connection_id.get(),
785 corr,
786 frame_type = ?frame.header.ty,
787 "routing control frame"
788 );
789 let dispatch_started_at = (frame.header.ty == FrameType::Request)
794 .then(|| dispatch_started_at.unwrap_or_else(Instant::now));
795 let responses = self
796 .control
797 .handle_control_frame_timed(ctx, frame, dispatch_started_at)
798 .await?;
799 for response in responses {
800 ctx.egress.send(response).await?;
801 }
802 return Ok(());
803 }
804
805 let data_route = self
806 .forwarding
807 .lookup_data_route(ctx.connection_id, channel, epoch)
808 .map_err(RouterError::Forwarding)?;
809
810 match data_route {
811 DataRoute::Module(DataRouteState::EpochMismatch) => {
812 if frame.header.ty == FrameType::Request {
813 self.counters
814 .increment_module_requests_dropped_stale_route();
815 let err = RouterError::StaleRouteEpoch {
816 channel,
817 epoch,
818 corr,
819 };
820 if let Some(error_frame) = err.to_error_frame() {
821 ctx.egress.send(error_frame).await?;
822 }
823 } else {
824 self.handle_orphan_module_frame(ctx, &frame)?;
825 }
826 debug!(
827 connection_id = ctx.connection_id.get(),
828 channel, epoch, corr, "dropping module frame for stale route epoch"
829 );
830 return Ok(());
831 }
832 DataRoute::Module(DataRouteState::Reserved) => {
833 if frame.header.ty == FrameType::Request {
834 self.counters
835 .increment_module_requests_dropped_stale_route();
836 let err = RouterError::UnknownChannel {
837 channel,
838 epoch,
839 corr,
840 };
841 if let Some(error_frame) = err.to_error_frame() {
842 ctx.egress.send(error_frame).await?;
843 }
844 } else {
845 self.record_module_frame_drop(ctx.connection_id)?;
846 }
847 debug!(
848 connection_id = ctx.connection_id.get(),
849 channel, epoch, corr, "dropping module frame for reserved route handle"
850 );
851 return Ok(());
852 }
853 DataRoute::Module(DataRouteState::Absent) => {
854 if frame.header.ty == FrameType::Request {
855 self.counters
856 .increment_module_requests_dropped_stale_route();
857 let err = RouterError::UnknownChannel {
858 channel,
859 epoch,
860 corr,
861 };
862 if let Some(error_frame) = err.to_error_frame() {
863 ctx.egress.send(error_frame).await?;
864 }
865 } else {
866 self.handle_orphan_module_frame(ctx, &frame)?;
867 }
868 debug!(
869 connection_id = ctx.connection_id.get(),
870 channel, epoch, corr, "dropping module frame for absent route handle"
871 );
872 return Ok(());
873 }
874 DataRoute::Module(DataRouteState::Bound(route)) => {
875 if frame.header.ty == FrameType::Goodbye {
876 if let RouteRelease::Removed(target) = self
877 .forwarding
878 .release_module_route(ctx.connection_id, channel, epoch)
879 .map_err(RouterError::Forwarding)?
880 {
881 let mut goodbye = frame;
882 goodbye.header.channel = target.channel;
883 goodbye.header.epoch = target.epoch;
884 if let Err(err) = target.sink.try_send(goodbye) {
885 if target.close_on_delivery_failure()
886 && self
887 .forwarding
888 .escalate_client_delivery_failure(
889 target.connection_id,
890 target.channel,
891 target.epoch,
892 CloseReason::new(
893 "route_goodbye_delivery_failed",
894 format!(
895 "failed to enqueue route GOODBYE for client channel {}: {err}",
896 target.channel
897 ),
898 ),
899 UndeliveredFrame {
900 module_id: Some(&route.module_id),
901 sink: &target.sink,
902 },
903 )
904 .map_err(RouterError::Forwarding)?
905 {
906 self.counters.increment_goodbye_relay_client_failed();
907 }
908 }
909 }
910 return Ok(());
911 }
912
913 let releases_credit = is_terminal_frame(frame.header.ty);
919 if releases_credit {
920 route.flow.release_corr(corr);
921 }
922 let mut frame = frame;
923 frame.header.channel = route.client_channel;
924 frame.header.epoch = route.client_epoch;
925 if let Err(err) = route.client_sink.try_send(frame) {
926 if self
927 .forwarding
928 .escalate_client_delivery_failure(
929 route.client_connection_id,
930 route.client_channel,
931 route.client_epoch,
932 CloseReason::new(
933 "module_to_client_delivery_failed",
934 format!(
935 "failed to enqueue module frame for client channel {} corr {corr}: {err}",
936 route.client_channel
937 ),
938 ),
939 UndeliveredFrame {
940 module_id: Some(&route.module_id),
941 sink: &route.client_sink,
942 },
943 )
944 .map_err(RouterError::Forwarding)?
945 {
946 self.counters
947 .increment_client_egress_close_delivery_failed();
948 }
949 return Ok(());
950 }
951 return Ok(());
952 }
953 DataRoute::Client(DataRouteState::EpochMismatch) => {
954 if frame.header.ty == FrameType::Request {
955 self.counters.increment_client_frames_dropped_stale_route();
956 let err = RouterError::StaleRouteEpoch {
958 channel,
959 epoch,
960 corr,
961 };
962 if let Some(error_frame) = err.to_error_frame() {
963 ctx.egress.send(error_frame).await?;
964 }
965 }
966 debug!(
967 connection_id = ctx.connection_id.get(),
968 channel, epoch, corr, "dropping client frame for stale route epoch"
969 );
970 return Ok(());
971 }
972 DataRoute::Client(DataRouteState::Reserved) => {
973 if frame.header.ty == FrameType::Request {
974 let err = RouterError::UnknownChannel {
975 channel,
976 epoch,
977 corr,
978 };
979 if let Some(error_frame) = err.to_error_frame() {
980 ctx.egress.send(error_frame).await?;
981 }
982 }
983 return Ok(());
984 }
985 DataRoute::Client(DataRouteState::Bound(route)) => {
986 if frame.header.ty == FrameType::Goodbye {
987 let _ = self
988 .control
989 .handle_route_goodbye(ctx.connection_id, channel, epoch)?;
990 return Ok(());
991 }
992 return self.forward_backend.handle_bound(frame, route).await;
993 }
994 DataRoute::Client(DataRouteState::Absent) => {}
995 }
996
997 if let Some(backend) = self.backends.get(&channel) {
998 return backend.handle(ctx.clone(), frame).await;
999 }
1000 if frame.header.ty == FrameType::Request {
1001 let err = RouterError::UnknownChannel {
1002 channel,
1003 epoch,
1004 corr,
1005 };
1006 if let Some(error_frame) = err.to_error_frame() {
1007 ctx.egress.send(error_frame).await?;
1008 }
1009 }
1010 Ok(())
1011 }
1012}
1013
1014impl Default for Router {
1015 fn default() -> Self {
1016 Self::with_default_self_handler()
1017 }
1018}
1019
1020#[must_use]
1022pub struct RouterConnection {
1023 id: ConnectionId,
1024 control_handler: Arc<ControlHandler>,
1025 forwarding: Arc<ForwardingTable>,
1026 close_receiver: Option<ConnectionCloseReceiver>,
1027 orphan_goodbyes: Arc<OrphanGoodbyeLimiter>,
1028}
1029
1030impl RouterConnection {
1031 pub fn id(&self) -> ConnectionId {
1032 self.id
1033 }
1034
1035 pub(crate) fn take_close_receiver(&mut self) -> ConnectionCloseReceiver {
1036 self.close_receiver
1037 .take()
1038 .expect("connection close receiver can only be taken once")
1039 }
1040}
1041
1042impl Drop for RouterConnection {
1043 fn drop(&mut self) {
1044 self.forwarding.unregister_connection_close(self.id);
1045 self.orphan_goodbyes.forget_connection(self.id);
1046 let _ = self.control_handler.cleanup_connection(self.id);
1049 }
1050}
1051
1052#[derive(Debug, Default, Clone, Copy)]
1055pub struct EchoBackend;
1056
1057impl EchoBackend {
1058 pub async fn handle(&self, ctx: RouteCtx, frame: Frame) -> Result<(), RouterError> {
1059 let response = Frame::build_with_version(
1060 frame.header.ver,
1061 FrameType::Response,
1062 frame.header.flags,
1063 frame.header.channel,
1064 frame.header.epoch,
1065 frame.header.corr,
1066 frame.body,
1067 )
1068 .map_err(RouterError::FrameBuild)?;
1069 ctx.egress.send(response).await
1070 }
1071}
1072
1073#[derive(Debug, Clone)]
1075pub struct ForwardBackend {
1076 forwarding: Arc<ForwardingTable>,
1077}
1078
1079impl ForwardBackend {
1080 pub fn new(forwarding: Arc<ForwardingTable>) -> Self {
1081 Self { forwarding }
1082 }
1083
1084 pub async fn handle(&self, ctx: RouteCtx, frame: Frame) -> Result<(), RouterError> {
1085 let channel = frame.header.channel;
1086 let corr = frame.header.corr;
1087 let route = match self
1088 .forwarding
1089 .lookup_data_route(ctx.connection_id, channel, frame.header.epoch)
1090 .map_err(RouterError::Forwarding)?
1091 {
1092 DataRoute::Client(DataRouteState::Bound(route)) => route,
1093 DataRoute::Client(_) | DataRoute::Module(_) => {
1094 return Err(RouterError::UnknownChannel {
1095 channel,
1096 epoch: frame.header.epoch,
1097 corr,
1098 });
1099 }
1100 };
1101 self.handle_bound(frame, route).await
1102 }
1103
1104 pub(crate) async fn handle_bound(
1105 &self,
1106 frame: Frame,
1107 route: Arc<RouteBinding>,
1108 ) -> Result<(), RouterError> {
1109 let channel = frame.header.channel;
1110 let corr = frame.header.corr;
1111 let frame_type = frame.header.ty;
1112
1113 let acquired_credit = frame_type == FrameType::Request;
1116 if acquired_credit {
1117 if let Err(err) = route
1118 .flow
1119 .acquire_tagged(corr, frame.header.flags.is_subscription())
1120 .await
1121 {
1122 if self
1129 .forwarding
1130 .endpoint_is_draining(route.module_endpoint)
1131 .map_err(RouterError::Forwarding)?
1132 {
1133 return Err(RouterError::route_error_with_epoch(
1134 channel,
1135 frame.header.epoch,
1136 corr,
1137 "module_reloading",
1138 format!("module endpoint for route channel {channel} is reloading"),
1139 ));
1140 }
1141 return Err(RouterError::backend_with_epoch(
1142 channel,
1143 frame.header.epoch,
1144 corr,
1145 format!("{err} for route channel {channel}"),
1146 ));
1147 }
1148 }
1149
1150 let mut frame = frame;
1151 frame.header.channel = route.module_channel;
1152 frame.header.epoch = route.module_epoch;
1153 let result = route.module_sink.send(frame).await.map_err(|err| {
1154 RouterError::backend_with_epoch(channel, route.client_epoch, corr, err.to_string())
1155 });
1156 if acquired_credit && result.is_err() {
1157 route.flow.release_corr(corr);
1158 }
1159 result
1160 }
1161}
1162
1163fn is_terminal_frame(frame_type: FrameType) -> bool {
1164 matches!(
1165 frame_type,
1166 FrameType::Response | FrameType::Error | FrameType::StreamEnd
1167 )
1168}
1169
1170#[derive(Debug, Clone, PartialEq, Eq)]
1173pub enum RouterError {
1174 ReservedChannelZero,
1175 DuplicateChannel {
1176 channel: u16,
1177 },
1178 UnknownChannel {
1179 channel: u16,
1180 epoch: u32,
1181 corr: u64,
1182 },
1183 StaleRouteEpoch {
1184 channel: u16,
1185 epoch: u32,
1186 corr: u64,
1187 },
1188 Backend {
1189 channel: u16,
1190 epoch: u32,
1191 corr: u64,
1192 message: String,
1193 },
1194 RouteError {
1195 channel: u16,
1196 epoch: u32,
1197 corr: u64,
1198 code: String,
1199 message: String,
1200 },
1201 FrameBuild(FrameBuildError),
1202 Forwarding(ForwardingError),
1203}
1204
1205impl RouterError {
1206 pub fn backend(channel: u16, corr: u64, message: impl Into<String>) -> Self {
1207 Self::backend_with_epoch(channel, 0, corr, message)
1208 }
1209
1210 pub fn backend_with_epoch(
1211 channel: u16,
1212 epoch: u32,
1213 corr: u64,
1214 message: impl Into<String>,
1215 ) -> Self {
1216 Self::Backend {
1217 channel,
1218 epoch,
1219 corr,
1220 message: message.into(),
1221 }
1222 }
1223
1224 pub fn route_error(
1225 channel: u16,
1226 corr: u64,
1227 code: impl Into<String>,
1228 message: impl Into<String>,
1229 ) -> Self {
1230 Self::route_error_with_epoch(channel, 0, corr, code, message)
1231 }
1232
1233 pub fn route_error_with_epoch(
1234 channel: u16,
1235 epoch: u32,
1236 corr: u64,
1237 code: impl Into<String>,
1238 message: impl Into<String>,
1239 ) -> Self {
1240 Self::RouteError {
1241 channel,
1242 epoch,
1243 corr,
1244 code: code.into(),
1245 message: message.into(),
1246 }
1247 }
1248
1249 pub fn to_error_frame(&self) -> Option<Frame> {
1251 match self {
1252 Self::UnknownChannel {
1253 channel,
1254 epoch,
1255 corr,
1256 } => error_frame(
1257 *channel,
1258 *epoch,
1259 *corr,
1260 error_codes::UNKNOWN_CHANNEL,
1261 format!("unknown channel {channel}"),
1262 ),
1263 Self::StaleRouteEpoch {
1264 channel,
1265 epoch,
1266 corr,
1267 } => error_frame(
1268 *channel,
1269 *epoch,
1270 *corr,
1271 error_codes::STALE_ROUTE_EPOCH,
1272 format!("stale route epoch for channel {channel}"),
1273 ),
1274 Self::Backend {
1275 channel,
1276 epoch,
1277 corr,
1278 message,
1279 } => error_frame(*channel, *epoch, *corr, "backend_error", message.clone()),
1280 Self::RouteError {
1281 channel,
1282 epoch,
1283 corr,
1284 code,
1285 message,
1286 } => error_frame(*channel, *epoch, *corr, code, message.clone()),
1287 Self::ReservedChannelZero
1288 | Self::DuplicateChannel { .. }
1289 | Self::FrameBuild(_)
1290 | Self::Forwarding(_) => None,
1291 }
1292 }
1293}
1294
1295fn error_frame(channel: u16, epoch: u32, corr: u64, code: &str, message: String) -> Option<Frame> {
1296 let body = serde_json::to_vec(&ErrorBody {
1297 code: code.to_string(),
1298 message,
1299 detail: None,
1300 })
1301 .ok()?;
1302
1303 Frame::build(
1304 FrameType::Error,
1305 Flags::new(false, Priority::Passive, false),
1306 channel,
1307 epoch,
1308 corr,
1309 body,
1310 )
1311 .ok()
1312}
1313
1314impl fmt::Display for RouterError {
1315 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1316 match self {
1317 Self::ReservedChannelZero => write!(f, "channel 0 is reserved for subc"),
1318 Self::DuplicateChannel { channel } => {
1319 write!(f, "backend already registered for channel {channel}")
1320 }
1321 Self::UnknownChannel { channel, corr, .. } => {
1322 write!(f, "unknown channel {channel} for corr {corr}")
1323 }
1324 Self::StaleRouteEpoch { channel, corr, .. } => {
1325 write!(f, "stale route epoch for channel {channel} corr {corr}")
1326 }
1327 Self::Backend {
1328 channel,
1329 corr,
1330 message,
1331 ..
1332 } => write!(
1333 f,
1334 "backend error on channel {channel} corr {corr}: {message}"
1335 ),
1336 Self::RouteError {
1337 channel,
1338 corr,
1339 code,
1340 message,
1341 ..
1342 } => write!(
1343 f,
1344 "route error {code} on channel {channel} corr {corr}: {message}"
1345 ),
1346 Self::FrameBuild(err) => write!(f, "failed to build routed frame: {err}"),
1347 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
1348 }
1349 }
1350}
1351
1352impl Error for RouterError {
1353 fn source(&self) -> Option<&(dyn Error + 'static)> {
1354 match self {
1355 Self::FrameBuild(err) => Some(err),
1356 Self::Forwarding(err) => Some(err),
1357 Self::ReservedChannelZero
1358 | Self::DuplicateChannel { .. }
1359 | Self::UnknownChannel { .. }
1360 | Self::StaleRouteEpoch { .. }
1361 | Self::Backend { .. }
1362 | Self::RouteError { .. } => None,
1363 }
1364 }
1365}
1366
1367#[cfg(test)]
1368mod tests {
1369 use super::*;
1370 use crate::{
1371 forwarding::RouteBindRelayOutcome,
1372 supervise::{ModuleSpec, RestartPolicy, Supervisor, SupervisorHandle},
1373 ControlHandler, Registry,
1374 };
1375 use std::{
1376 sync::{mpsc as std_mpsc, Arc},
1377 time::Duration,
1378 };
1379 use subc_control::ModuleProtocol;
1380 use subc_protocol::{manifest::Concurrency, ErrorBody, Flags, FrameType, Priority};
1381 use tokio::sync::mpsc;
1382
1383 pub(crate) use crate::router::test_log::{captured_logs, log_capture};
1384
1385 fn logged_millis(logs: &str, field: &str) -> u64 {
1386 logs.split_whitespace()
1387 .find_map(|part| part.strip_prefix(field))
1388 .and_then(|value| value.parse().ok())
1389 .unwrap_or_else(|| panic!("missing numeric {field} in logs: {logs}"))
1390 }
1391
1392 fn request(channel: u16, corr: u64, body: &[u8]) -> Frame {
1393 Frame::build(
1394 FrameType::Request,
1395 Flags::new(true, Priority::Interactive, false),
1396 channel,
1397 0,
1398 corr,
1399 body.to_vec(),
1400 )
1401 .unwrap()
1402 }
1403
1404 fn ping(corr: u64) -> Frame {
1405 Frame::build(
1406 FrameType::Ping,
1407 Flags::new(false, Priority::Passive, false),
1408 0,
1409 0,
1410 corr,
1411 Vec::new(),
1412 )
1413 .unwrap()
1414 }
1415
1416 fn route_ctx() -> (RouteCtx, mpsc::Receiver<crate::router::OutboundFrame>) {
1417 let (tx, rx) = mpsc::channel(8);
1418 (
1419 RouteCtx {
1420 connection_id: ConnectionId::LOCAL,
1421 egress: FrameSink::new(tx),
1422 },
1423 rx,
1424 )
1425 }
1426
1427 #[tokio::test]
1428 async fn echo_backend_returns_response_with_byte_identical_body() {
1429 let mut router = Router::with_default_self_handler();
1430 router.register_backend(7, EchoBackend).unwrap();
1431 let (ctx, mut rx) = route_ctx();
1432 let body = b"{not parsed}\0\xff";
1433
1434 router
1435 .route_for_connection(&ctx, request(7, 123, body))
1436 .await
1437 .unwrap();
1438 let response = rx.recv().await.unwrap();
1439
1440 assert_eq!(response.header.ty, FrameType::Response);
1441 assert_eq!(response.header.channel, 7);
1442 assert_eq!(response.header.corr, 123);
1443 assert_eq!(response.body, body);
1444 assert!(rx.try_recv().is_err());
1445 }
1446
1447 #[tokio::test]
1448 async fn unknown_channel_emits_canonical_error_frame() {
1449 let router = Router::with_default_self_handler();
1450 let (ctx, mut rx) = route_ctx();
1451
1452 router
1453 .route_for_connection(&ctx, request(99, 5, b"payload"))
1454 .await
1455 .unwrap();
1456 let error_frame = rx.recv().await.unwrap();
1457
1458 assert_eq!(error_frame.header.ty, FrameType::Error);
1459 assert_eq!(error_frame.header.channel, 99);
1460 assert_eq!(error_frame.header.corr, 5);
1461 let body: ErrorBody = serde_json::from_slice(&error_frame.body).unwrap();
1462 assert_eq!(body.code, "unknown_channel");
1463 assert_eq!(body.message, "unknown channel 99");
1464 }
1465
1466 #[tokio::test]
1467 async fn channel_zero_uses_control_handler_not_backend_registry() {
1468 let mut router = Router::with_default_self_handler();
1469 router.register_backend(1, EchoBackend).unwrap();
1470 let (ctx, mut rx) = route_ctx();
1471
1472 router.route_for_connection(&ctx, ping(77)).await.unwrap();
1473 let response = rx.recv().await.unwrap();
1474
1475 assert_eq!(response.header.ty, FrameType::Pong);
1476 assert_eq!(response.header.channel, 0);
1477 assert_eq!(response.header.corr, 77);
1478 assert!(response.body.is_empty());
1479 }
1480
1481 #[tokio::test]
1482 async fn slow_control_dispatch_logs_decoded_op_and_elapsed_time() {
1483 let control = Arc::new(
1484 ControlHandler::new(Arc::new(Registry::default()))
1485 .with_control_dispatch_delay(Duration::from_millis(1050)),
1486 );
1487 let router = Router::with_control_handler(control);
1488 let (ctx, mut rx) = route_ctx();
1489 let (output, guard) = log_capture(tracing::Level::WARN);
1490
1491 router
1492 .route_for_connection(&ctx, request(0, 41, br#"{"op":"server.describe"}"#))
1493 .await
1494 .expect("slow request routes");
1495 assert!(rx.recv().await.is_some(), "request receives a response");
1496 drop(guard);
1497
1498 let logs = captured_logs(&output);
1499 assert!(logs.contains("slow control dispatch"));
1500 assert!(logs.contains("op=server.describe"));
1501 assert!(logs.contains("connection_id=0"));
1502 assert!(logs.contains("corr=41"));
1503 assert!(
1504 logged_millis(&logs, "elapsed_ms=") >= 1050,
1505 "elapsed must include the injected handler delay: {logs}"
1506 );
1507 }
1508
1509 #[tokio::test]
1510 async fn fast_control_dispatch_emits_arrival_without_slow_warning() {
1511 let router = Router::with_default_self_handler();
1512 let (ctx, mut rx) = route_ctx();
1513 let (output, guard) = log_capture(tracing::Level::DEBUG);
1514
1515 router
1516 .route_for_connection(&ctx, request(0, 42, br#"{"op":"server.describe"}"#))
1517 .await
1518 .expect("fast request routes");
1519 assert!(rx.recv().await.is_some(), "request receives a response");
1520 drop(guard);
1521
1522 let logs = captured_logs(&output);
1523 assert!(logs.contains("control dispatch op=server.describe connection_id=0 corr=42"));
1524 assert!(!logs.contains("slow control dispatch"));
1525 }
1526
1527 #[tokio::test]
1528 async fn control_dispatch_arrival_is_hidden_at_info() {
1529 let router = Router::with_default_self_handler();
1530 let (ctx, mut rx) = route_ctx();
1531 let (output, guard) = log_capture(tracing::Level::INFO);
1532
1533 router
1534 .route_for_connection(&ctx, request(0, 43, br#"{"op":"server.describe"}"#))
1535 .await
1536 .expect("fast request routes");
1537 assert!(rx.recv().await.is_some(), "request receives a response");
1538 drop(guard);
1539
1540 assert!(
1541 !captured_logs(&output).contains("control dispatch"),
1542 "arrival logging must stay hidden at INFO"
1543 );
1544 }
1545
1546 #[tokio::test]
1547 async fn supervisor_list_logs_contended_snapshot_lock_only() {
1548 let registry = Arc::new(Registry::default());
1549 let handle = SupervisorHandle::new();
1550 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
1551 .with_handle(handle.clone());
1552 let module = supervisor
1553 .supervise_configured(
1554 ModuleSpec {
1555 launch_nonce_env: true,
1556 module_id: "held-module".to_string(),
1557 program: "test-module".into(),
1558 args: Vec::new(),
1559 env: Vec::new(),
1560 reserved: false,
1561 reserved_prefixes: Vec::new(),
1562 protocol: ModuleProtocol::Subc,
1563 overlap: Default::default(),
1564 },
1565 false,
1566 )
1567 .expect("disabled test module is supervised");
1568 let router = Router::with_control_handler(Arc::new(
1569 ControlHandler::new(Arc::clone(®istry)).with_supervisor(handle),
1570 ));
1571 let (ctx, mut rx) = route_ctx();
1572 let (acquired, ready) = std_mpsc::channel();
1573 let holder = module.hold_snapshot_for_test(acquired, Duration::from_millis(400));
1574 ready.recv().expect("holder acquired snapshot lock");
1575 let (output, guard) = log_capture(tracing::Level::WARN);
1576
1577 router
1578 .route_for_connection(&ctx, request(0, 44, br#"{"op":"supervisor.list"}"#))
1579 .await
1580 .expect("list request routes after the lock releases");
1581 assert!(
1582 rx.recv().await.is_some(),
1583 "list request receives a response"
1584 );
1585 holder.join().expect("snapshot holder exits cleanly");
1586 drop(guard);
1587
1588 let logs = captured_logs(&output);
1589 assert!(logs.contains("slow snapshot lock"));
1590 assert!(logs.contains("module_id=held-module"));
1591 assert!(logs.contains("caller=list"));
1592 assert!(
1593 logged_millis(&logs, "waited_ms=") >= 250,
1594 "wait must exceed the slow-lock threshold: {logs}"
1595 );
1596
1597 let (output, guard) = log_capture(tracing::Level::WARN);
1598 router
1599 .route_for_connection(&ctx, request(0, 45, br#"{"op":"supervisor.list"}"#))
1600 .await
1601 .expect("uncontended list request routes");
1602 assert!(
1603 rx.recv().await.is_some(),
1604 "uncontended list receives a response"
1605 );
1606 drop(guard);
1607 assert!(
1608 !captured_logs(&output).contains("slow snapshot lock"),
1609 "uncontended list acquisition must not warn"
1610 );
1611 }
1612
1613 #[tokio::test]
1614 async fn full_module_to_client_sink_requests_client_close_without_erroring_module() {
1615 let forwarding = Arc::new(ForwardingTable::default());
1616 let control = Arc::new(ControlHandler::with_forwarding(
1617 Arc::new(crate::Registry::default()),
1618 Arc::clone(&forwarding),
1619 ));
1620 let router = Router::with_control_handler(control);
1621 let module_connection = ConnectionId::new(10);
1622 let client_connection = ConnectionId::new(20);
1623 let mut close_receiver = forwarding.register_connection_close(client_connection);
1624 let (module_tx, _module_rx) = mpsc::channel(1);
1625 forwarding
1626 .register_module_connection(
1627 module_connection,
1628 "full-sink-provider".to_string(),
1629 1,
1630 Concurrency::ModuleManaged,
1631 FrameSink::new(module_tx),
1632 )
1633 .unwrap();
1634 let (client_tx, mut client_rx) = mpsc::channel(1);
1635 let pending = forwarding
1636 .begin_route_bind_relay_for_test(
1637 client_connection,
1638 FrameSink::new(client_tx),
1639 700,
1640 "full-sink-provider",
1641 )
1642 .unwrap();
1643 forwarding
1644 .complete_pending_relay(
1645 module_connection,
1646 pending.corr,
1647 RouteBindRelayOutcome::Accepted,
1648 )
1649 .unwrap();
1650
1651 let (module_egress_tx, _module_egress_rx) = mpsc::channel(1);
1652 let module_ctx = RouteCtx {
1653 connection_id: module_connection,
1654 egress: FrameSink::new(module_egress_tx),
1655 };
1656 let terminal = Frame::build(
1657 FrameType::Response,
1658 Flags::new(false, Priority::Interactive, true),
1659 pending.module_channel,
1660 pending.module_epoch,
1661 701,
1662 b"terminal".to_vec(),
1663 )
1664 .unwrap();
1665
1666 router
1667 .route_for_connection(&module_ctx, terminal)
1668 .await
1669 .unwrap();
1670 let reason = tokio::time::timeout(Duration::from_secs(1), &mut close_receiver)
1671 .await
1672 .expect("close request should be sent for the full client sink")
1673 .expect("close sender should include a reason");
1674 assert!(
1675 reason
1676 .to_string()
1677 .contains("module_to_client_delivery_failed"),
1678 "unexpected close reason: {reason}"
1679 );
1680 assert_eq!(client_rx.try_recv().unwrap().header.corr, 700);
1681 assert!(client_rx.try_recv().is_err());
1682 assert_eq!(
1683 router.counters.snapshot()["client_egress_close_delivery_failed"],
1684 1
1685 );
1686 }
1687
1688 #[tokio::test]
1692 async fn terminal_frame_releases_its_credit_even_when_client_delivery_fails() {
1693 let forwarding = Arc::new(ForwardingTable::default());
1694 let control = Arc::new(ControlHandler::with_forwarding(
1695 Arc::new(crate::Registry::default()),
1696 Arc::clone(&forwarding),
1697 ));
1698 let router = Router::with_control_handler(control);
1699 let module_connection = ConnectionId::new(11);
1700 let client_connection = ConnectionId::new(21);
1701 let _close_receiver = forwarding.register_connection_close(client_connection);
1702 let (module_tx, _module_rx) = mpsc::channel(1);
1703 forwarding
1704 .register_module_connection(
1705 module_connection,
1706 "credit-provider".to_string(),
1707 1,
1708 Concurrency::ModuleManaged,
1709 FrameSink::new(module_tx),
1710 )
1711 .unwrap();
1712 let (client_tx, _client_rx) = mpsc::channel(1);
1715 let pending = forwarding
1716 .begin_route_bind_relay_for_test(
1717 client_connection,
1718 FrameSink::new(client_tx),
1719 800,
1720 "credit-provider",
1721 )
1722 .unwrap();
1723 forwarding
1724 .complete_pending_relay(
1725 module_connection,
1726 pending.corr,
1727 RouteBindRelayOutcome::Accepted,
1728 )
1729 .unwrap();
1730 let DataRoute::Client(DataRouteState::Bound(route)) = forwarding
1731 .lookup_data_route(
1732 client_connection,
1733 pending.client_channel,
1734 pending.client_epoch,
1735 )
1736 .unwrap()
1737 else {
1738 panic!("expected a bound client route");
1739 };
1740 route.flow.acquire_tagged(801, false).await.unwrap();
1741 assert_eq!(route.flow.drain_in_flight(), 1);
1742
1743 let (module_egress_tx, _module_egress_rx) = mpsc::channel(1);
1744 let module_ctx = RouteCtx {
1745 connection_id: module_connection,
1746 egress: FrameSink::new(module_egress_tx),
1747 };
1748 let terminal = Frame::build(
1749 FrameType::Response,
1750 Flags::new(false, Priority::Interactive, true),
1751 pending.module_channel,
1752 pending.module_epoch,
1753 801,
1754 b"terminal".to_vec(),
1755 )
1756 .unwrap();
1757 router
1758 .route_for_connection(&module_ctx, terminal)
1759 .await
1760 .unwrap();
1761
1762 assert_eq!(
1763 router.counters.snapshot()["client_egress_close_delivery_failed"],
1764 1,
1765 "the client delivery must have failed for this test to mean anything"
1766 );
1767 assert_eq!(
1768 route.flow.drain_in_flight(),
1769 0,
1770 "the module's terminal frame must release its credit even though the client could not take it"
1771 );
1772 }
1773
1774 #[tokio::test]
1775 async fn full_route_goodbye_sink_requests_target_close_without_erroring_module() {
1776 let forwarding = Arc::new(ForwardingTable::default());
1777 let control = Arc::new(ControlHandler::with_forwarding(
1778 Arc::new(crate::Registry::default()),
1779 Arc::clone(&forwarding),
1780 ));
1781 let router = Router::with_control_handler(control);
1782 let module_connection = ConnectionId::new(30);
1783 let client_connection = ConnectionId::new(40);
1784 let mut close_receiver = forwarding.register_connection_close(client_connection);
1785 let (module_tx, _module_rx) = mpsc::channel(1);
1786 forwarding
1787 .register_module_connection(
1788 module_connection,
1789 "goodbye-full-provider".to_string(),
1790 1,
1791 Concurrency::ModuleManaged,
1792 FrameSink::new(module_tx),
1793 )
1794 .unwrap();
1795 let (client_tx, mut client_rx) = mpsc::channel(1);
1796 let pending = forwarding
1797 .begin_route_bind_relay_for_test(
1798 client_connection,
1799 FrameSink::new(client_tx),
1800 800,
1801 "goodbye-full-provider",
1802 )
1803 .unwrap();
1804 forwarding
1805 .complete_pending_relay(
1806 module_connection,
1807 pending.corr,
1808 RouteBindRelayOutcome::Accepted,
1809 )
1810 .unwrap();
1811
1812 let (module_egress_tx, _module_egress_rx) = mpsc::channel(1);
1813 let module_ctx = RouteCtx {
1814 connection_id: module_connection,
1815 egress: FrameSink::new(module_egress_tx),
1816 };
1817 let goodbye = Frame::build(
1818 FrameType::Goodbye,
1819 Flags::new(false, Priority::Passive, true),
1820 pending.module_channel,
1821 pending.module_epoch,
1822 801,
1823 Vec::new(),
1824 )
1825 .unwrap();
1826
1827 router
1828 .route_for_connection(&module_ctx, goodbye)
1829 .await
1830 .unwrap();
1831 let reason = tokio::time::timeout(Duration::from_secs(1), &mut close_receiver)
1832 .await
1833 .expect("close request should be sent for the full GOODBYE sink")
1834 .expect("close sender should include a reason");
1835 assert!(
1836 reason.to_string().contains("route_goodbye_delivery_failed"),
1837 "unexpected close reason: {reason}"
1838 );
1839 assert_eq!(client_rx.try_recv().unwrap().header.corr, 800);
1840 assert!(client_rx.try_recv().is_err());
1841 assert_eq!(router.counters.snapshot()["goodbye_relay_client_failed"], 1);
1842 assert_eq!(router.counters.snapshot()["route_released_epoch_fenced"], 1);
1843 }
1844
1845 async fn multi_route_client(
1850 module_id: &str,
1851 module_connection: ConnectionId,
1852 client_connection: ConnectionId,
1853 routes: usize,
1854 ) -> (
1855 Router,
1856 FrameSink,
1857 mpsc::Receiver<OutboundFrame>,
1858 RouteCtx,
1859 Vec<crate::forwarding::PendingRouteBindRelay>,
1860 ConnectionCloseReceiver,
1861 ) {
1862 let forwarding = Arc::new(ForwardingTable::default());
1863 let control = Arc::new(ControlHandler::with_forwarding(
1864 Arc::new(crate::Registry::default()),
1865 Arc::clone(&forwarding),
1866 ));
1867 let router = Router::with_control_handler(control);
1868 let close_receiver = forwarding.register_connection_close(client_connection);
1869 let (module_tx, _module_rx) = mpsc::channel(8);
1870 forwarding
1871 .register_module_connection(
1872 module_connection,
1873 module_id.to_string(),
1874 1,
1875 Concurrency::ModuleManaged,
1876 FrameSink::new(module_tx),
1877 )
1878 .unwrap();
1879 let (client_sink, mut client_rx) = crate::server::connection_egress();
1880 let mut bound = Vec::with_capacity(routes);
1881 for index in 0..routes {
1882 let pending = forwarding
1883 .begin_route_bind_relay_for_test(
1884 client_connection,
1885 client_sink.clone(),
1886 900 + index as u64,
1887 module_id,
1888 )
1889 .unwrap();
1890 forwarding
1891 .complete_pending_relay(
1892 module_connection,
1893 pending.corr,
1894 RouteBindRelayOutcome::Accepted,
1895 )
1896 .unwrap();
1897 assert_eq!(
1898 client_rx.recv().await.unwrap().header.corr,
1899 900 + index as u64
1900 );
1901 bound.push(pending);
1902 }
1903 let (module_egress_tx, _module_egress_rx) = mpsc::channel(8);
1904 let module_ctx = RouteCtx {
1905 connection_id: module_connection,
1906 egress: FrameSink::new(module_egress_tx),
1907 };
1908 (
1909 router,
1910 client_sink,
1911 client_rx,
1912 module_ctx,
1913 bound,
1914 close_receiver,
1915 )
1916 }
1917
1918 #[test]
1925 #[ignore = "timing measurement, run on demand"]
1926 fn egress_sink_per_frame_cost() {
1927 const FRAMES: usize = 1_000_000;
1928 const BATCH: usize = 1_000;
1929 let (sink, mut rx) = crate::server::connection_egress();
1930 let template = stream_frame(9, 1, 0, vec![b't'; 200]);
1931 let started = Instant::now();
1932 for _ in 0..FRAMES / BATCH {
1933 for _ in 0..BATCH {
1934 sink.try_send(template.clone()).unwrap();
1935 }
1936 for _ in 0..BATCH {
1937 let outbound = rx.try_recv().unwrap();
1938 if let Some(charge) = &outbound.charge {
1939 charge.taken_by_writer();
1940 }
1941 drop(std::hint::black_box(outbound));
1942 }
1943 }
1944 let elapsed = started.elapsed();
1945 assert_eq!(sink.backlog().queued_bytes, 0);
1946 println!(
1947 "egress sink: {FRAMES} frames in {elapsed:?}, {:.1} ns/frame",
1948 elapsed.as_nanos() as f64 / FRAMES as f64
1949 );
1950 }
1951
1952 fn stream_frame(channel: u16, epoch: u32, corr: u64, body: Vec<u8>) -> Frame {
1953 Frame::build(
1954 FrameType::StreamData,
1955 Flags::new(false, Priority::Interactive, false),
1956 channel,
1957 epoch,
1958 corr,
1959 body,
1960 )
1961 .unwrap()
1962 }
1963
1964 #[tokio::test]
1970 async fn awaited_send_parked_behind_a_full_byte_budget_wakes_when_bytes_free() {
1971 let (tx, mut rx) = mpsc::channel(1024);
1972 let sink = FrameSink::with_byte_budget(tx, 2_000);
1973 let mut queued = 0u64;
1974 while sink
1975 .try_send(stream_frame(7, 1, queued, vec![b'x'; 200]))
1976 .is_ok()
1977 {
1978 queued += 1;
1979 }
1980 assert!(queued > 0, "the budget admitted nothing");
1981
1982 let parked = tokio::spawn({
1983 let sink = sink.clone();
1984 async move { sink.send(stream_frame(0, 0, 999, vec![b'r'; 200])).await }
1985 });
1986 tokio::time::sleep(Duration::from_millis(100)).await;
1987 assert!(
1988 !parked.is_finished(),
1989 "the awaited send must wait while the byte budget is full"
1990 );
1991
1992 drop(rx.recv().await.expect("a queued frame"));
1994 tokio::time::timeout(Duration::from_secs(2), parked)
1995 .await
1996 .expect("the parked send was never woken after bytes were freed")
1997 .unwrap()
1998 .unwrap();
1999 }
2000
2001 #[tokio::test]
2006 async fn paused_client_reader_keeps_connection_through_small_frame_burst() {
2007 const ROUTES: usize = 4;
2008 const FRAMES_PER_ROUTE: usize = 250;
2009 let (router, client_sink, mut client_rx, module_ctx, routes, mut close_receiver) =
2010 multi_route_client(
2011 "burst-provider",
2012 ConnectionId::new(60),
2013 ConnectionId::new(61),
2014 ROUTES,
2015 )
2016 .await;
2017
2018 for seq in 0..FRAMES_PER_ROUTE as u64 {
2020 for (index, route) in routes.iter().enumerate() {
2021 let body = format!("route-{index}-token-{seq:05}-{}", "t".repeat(170));
2022 router
2023 .route_for_connection(
2024 &module_ctx,
2025 stream_frame(
2026 route.module_channel,
2027 route.module_epoch,
2028 seq,
2029 body.into_bytes(),
2030 ),
2031 )
2032 .await
2033 .unwrap();
2034 }
2035 }
2036
2037 let backlog = client_sink.backlog();
2038 assert_eq!(backlog.queued_frames, ROUTES * FRAMES_PER_ROUTE);
2039 assert!(backlog.queued_bytes < crate::server::CONNECTION_EGRESS_BYTE_BUDGET);
2040 assert!(
2041 close_receiver.try_recv().is_err(),
2042 "a paused reader under the byte budget must not be closed"
2043 );
2044 assert_eq!(
2045 router.counters.snapshot()["client_egress_close_delivery_failed"],
2046 0
2047 );
2048
2049 let mut next_seq = vec![0u64; ROUTES];
2051 for _ in 0..ROUTES * FRAMES_PER_ROUTE {
2052 let frame = client_rx.try_recv().expect("every queued frame arrives");
2053 let index = routes
2054 .iter()
2055 .position(|route| route.client_channel == frame.header.channel)
2056 .expect("frame arrives on one of the bound client channels");
2057 assert_eq!(frame.header.corr, next_seq[index], "route {index} order");
2058 let expected_prefix = format!("route-{index}-token-{:05}-", next_seq[index]);
2059 assert!(frame.body.starts_with(expected_prefix.as_bytes()));
2060 next_seq[index] += 1;
2061 }
2062 assert!(client_rx.try_recv().is_err());
2063 assert_eq!(next_seq, vec![FRAMES_PER_ROUTE as u64; ROUTES]);
2064 assert_eq!(client_sink.backlog().queued_bytes, 0);
2065 }
2066
2067 #[tokio::test]
2071 async fn never_reading_client_is_closed_at_byte_budget_with_warn_diagnosis() {
2072 let (logs, _guard) = test_log::log_capture(tracing::Level::WARN);
2073 const BODY: usize = 16 * 1024;
2074 let (router, client_sink, _client_rx, module_ctx, routes, mut close_receiver) =
2075 multi_route_client(
2076 "stuck-reader-provider",
2077 ConnectionId::new(70),
2078 ConnectionId::new(71),
2079 2,
2080 )
2081 .await;
2082
2083 let mut admitted = 0usize;
2084 let mut sent = 0u64;
2085 while router.counters.snapshot()["client_egress_close_delivery_failed"] == 0 {
2086 assert!(sent < 1_000, "the byte budget never refused a frame");
2087 tracing::callsite::rebuild_interest_cache();
2094 let route = &routes[(sent % 2) as usize];
2095 router
2096 .route_for_connection(
2097 &module_ctx,
2098 stream_frame(
2099 route.module_channel,
2100 route.module_epoch,
2101 sent,
2102 vec![b'z'; BODY],
2103 ),
2104 )
2105 .await
2106 .unwrap();
2107 sent += 1;
2108 admitted = client_sink.backlog().queued_frames;
2109 }
2110 let frame_bytes = subc_protocol::HEADER_LEN + BODY;
2112 assert_eq!(
2113 admitted,
2114 crate::server::CONNECTION_EGRESS_BYTE_BUDGET / frame_bytes
2115 );
2116 let reason = close_receiver
2117 .try_recv()
2118 .expect("the client connection must be asked to close");
2119 assert!(reason
2120 .to_string()
2121 .contains("module_to_client_delivery_failed"));
2122
2123 router
2125 .route_for_connection(
2126 &module_ctx,
2127 stream_frame(
2128 routes[0].module_channel,
2129 routes[0].module_epoch,
2130 sent,
2131 vec![b'z'; BODY],
2132 ),
2133 )
2134 .await
2135 .unwrap();
2136
2137 let captured = test_log::captured_logs(&logs);
2138 let warn_lines = captured
2139 .lines()
2140 .filter(|line| {
2141 line.contains("closing client connection: its egress queue could not take a frame")
2142 })
2143 .collect::<Vec<_>>();
2144 assert_eq!(warn_lines.len(), 1, "exactly one WARN, got: {captured}");
2145 let line = warn_lines[0];
2146 assert!(line.contains("WARN"), "{line}");
2147 assert!(line.contains("connection_id=71"), "{line}");
2148 assert!(
2149 line.contains("module_id=\"stuck-reader-provider\""),
2150 "{line}"
2151 );
2152 assert!(line.contains("client_channel="), "{line}");
2153 assert!(line.contains("principals=direct"), "{line}");
2154 let queued_bytes: usize = line
2155 .split("queued_bytes=")
2156 .nth(1)
2157 .and_then(|rest| rest.split_whitespace().next())
2158 .and_then(|value| value.parse().ok())
2159 .expect("queued_bytes is logged");
2160 assert_eq!(queued_bytes, admitted * frame_bytes);
2161 assert!(
2162 line.contains(&format!("queued_frames={admitted}")),
2163 "{line}"
2164 );
2165 assert!(line.contains("oldest_queued_ms="), "{line}");
2166 }
2167
2168 fn route_frame(ty: FrameType, channel: u16, epoch: u32, corr: u64) -> Frame {
2169 Frame::build(
2170 ty,
2171 Flags::new(false, Priority::Interactive, false),
2172 channel,
2173 epoch,
2174 corr,
2175 if ty == FrameType::Request || ty == FrameType::Response {
2176 b"route-body".to_vec()
2177 } else {
2178 Vec::new()
2179 },
2180 )
2181 .unwrap()
2182 }
2183
2184 type DynamicRouteFixture = (
2185 Router,
2186 Arc<ForwardingTable>,
2187 RouteCtx,
2188 mpsc::Receiver<crate::router::OutboundFrame>,
2189 RouteCtx,
2190 mpsc::Receiver<crate::router::OutboundFrame>,
2191 mpsc::Receiver<crate::router::OutboundFrame>,
2192 crate::forwarding::PendingRouteBindRelay,
2193 );
2194
2195 fn dynamic_route_fixture(commit: bool) -> DynamicRouteFixture {
2196 let forwarding = Arc::new(ForwardingTable::default());
2197 let control = Arc::new(crate::ControlHandler::with_forwarding(
2198 Arc::new(crate::Registry::default()),
2199 Arc::clone(&forwarding),
2200 ));
2201 let router = Router::with_control_handler(control);
2202 let module_connection = ConnectionId::new(500);
2203 let client_connection = ConnectionId::new(501);
2204 let (module_tx, module_rx) = mpsc::channel(8);
2205 forwarding
2206 .register_module_connection(
2207 module_connection,
2208 "epoch-router".into(),
2209 2,
2210 Concurrency::ModuleManaged,
2211 FrameSink::new(module_tx),
2212 )
2213 .unwrap();
2214 let (client_tx, client_rx) = mpsc::channel(8);
2215 let client_sink = FrameSink::new(client_tx);
2216 let pending = forwarding
2217 .begin_route_bind_relay_for_test(
2218 client_connection,
2219 client_sink.clone(),
2220 700,
2221 "epoch-router",
2222 )
2223 .unwrap();
2224 if commit {
2225 forwarding
2226 .complete_pending_relay(
2227 module_connection,
2228 pending.corr,
2229 RouteBindRelayOutcome::Accepted,
2230 )
2231 .unwrap();
2232 }
2233 let (module_egress_tx, module_egress_rx) = mpsc::channel(8);
2234 (
2235 router,
2236 forwarding,
2237 RouteCtx {
2238 connection_id: client_connection,
2239 egress: client_sink,
2240 },
2241 client_rx,
2242 RouteCtx {
2243 connection_id: module_connection,
2244 egress: FrameSink::new(module_egress_tx),
2245 },
2246 module_egress_rx,
2247 module_rx,
2248 pending,
2249 )
2250 }
2251
2252 #[tokio::test]
2253 async fn route_epochs_validate_both_directions_and_rewrite_to_peer_handle() {
2254 let (
2255 router,
2256 _forwarding,
2257 client_ctx,
2258 mut client_rx,
2259 module_ctx,
2260 _module_egress_rx,
2261 mut module_rx,
2262 pending,
2263 ) = dynamic_route_fixture(true);
2264 let route_open = client_rx.recv().await.unwrap();
2265 assert_eq!(route_open.header.corr, 700);
2266
2267 router
2268 .route_for_connection(
2269 &client_ctx,
2270 route_frame(
2271 FrameType::Request,
2272 pending.client_channel,
2273 pending.client_epoch,
2274 701,
2275 ),
2276 )
2277 .await
2278 .unwrap();
2279 let forwarded = module_rx.recv().await.unwrap();
2280 assert_eq!(forwarded.header.channel, pending.module_channel);
2281 assert_eq!(forwarded.header.epoch, pending.module_epoch);
2282
2283 router
2284 .route_for_connection(
2285 &module_ctx,
2286 route_frame(
2287 FrameType::Response,
2288 pending.module_channel,
2289 pending.module_epoch,
2290 701,
2291 ),
2292 )
2293 .await
2294 .unwrap();
2295 let delivered = client_rx.recv().await.unwrap();
2296 assert_eq!(delivered.header.channel, pending.client_channel);
2297 assert_eq!(delivered.header.epoch, pending.client_epoch);
2298
2299 router
2300 .route_for_connection(
2301 &client_ctx,
2302 route_frame(
2303 FrameType::Request,
2304 pending.client_channel,
2305 pending.client_epoch + 1,
2306 702,
2307 ),
2308 )
2309 .await
2310 .unwrap();
2311 router
2312 .route_for_connection(
2313 &module_ctx,
2314 route_frame(
2315 FrameType::Response,
2316 pending.module_channel,
2317 pending.module_epoch + 1,
2318 703,
2319 ),
2320 )
2321 .await
2322 .unwrap();
2323 let stale_error = client_rx.recv().await.unwrap();
2324 assert_eq!(stale_error.header.ty, FrameType::Error);
2325 assert_eq!(stale_error.header.channel, pending.client_channel);
2326 assert_eq!(stale_error.header.epoch, pending.client_epoch + 1);
2327 assert_eq!(stale_error.header.corr, 702);
2328 let body: ErrorBody = serde_json::from_slice(&stale_error.body).unwrap();
2329 assert_eq!(body.code, "stale_route_epoch");
2330 assert!(module_rx.try_recv().is_err());
2331 assert!(client_rx.try_recv().is_err());
2332 let counters = router.counters.snapshot();
2333 assert_eq!(counters["client_frames_dropped_stale_route"], 1);
2334 assert_eq!(counters["module_frames_dropped_no_route"], 1);
2335 }
2336
2337 #[tokio::test]
2338 async fn accepted_route_publishes_route_open_before_immediate_reverse_request() {
2339 let (
2340 router,
2341 _,
2342 _client_ctx,
2343 mut client_rx,
2344 module_ctx,
2345 _module_egress_rx,
2346 _module_rx,
2347 pending,
2348 ) = dynamic_route_fixture(true);
2349 router
2350 .route_for_connection(
2351 &module_ctx,
2352 route_frame(
2353 FrameType::Request,
2354 pending.module_channel,
2355 pending.module_epoch,
2356 800,
2357 ),
2358 )
2359 .await
2360 .unwrap();
2361
2362 let first = client_rx.recv().await.unwrap();
2363 let second = client_rx.recv().await.unwrap();
2364 assert_eq!(first.header.channel, 0);
2365 assert_eq!(first.header.corr, 700);
2366 assert_eq!(second.header.channel, pending.client_channel);
2367 assert_eq!(second.header.epoch, pending.client_epoch);
2368 assert_eq!(second.header.corr, 800);
2369 }
2370
2371 #[tokio::test]
2372 async fn reserved_slot_ingress_errors_only_matching_client_requests() {
2373 let (
2374 router,
2375 _forwarding,
2376 client_ctx,
2377 mut client_rx,
2378 _module_ctx,
2379 _module_egress_rx,
2380 mut module_rx,
2381 pending,
2382 ) = dynamic_route_fixture(false);
2383 router
2384 .route_for_connection(
2385 &client_ctx,
2386 route_frame(
2387 FrameType::Request,
2388 pending.client_channel,
2389 pending.client_epoch,
2390 900,
2391 ),
2392 )
2393 .await
2394 .unwrap();
2395 let error = client_rx.recv().await.unwrap();
2396 assert_eq!(error.header.ty, FrameType::Error);
2397 assert_eq!(error.header.channel, pending.client_channel);
2398 assert_eq!(error.header.epoch, pending.client_epoch);
2399 assert_eq!(error.header.corr, 900);
2400
2401 router
2402 .route_for_connection(
2403 &client_ctx,
2404 route_frame(
2405 FrameType::Response,
2406 pending.client_channel,
2407 pending.client_epoch,
2408 901,
2409 ),
2410 )
2411 .await
2412 .unwrap();
2413 router
2414 .route_for_connection(
2415 &client_ctx,
2416 route_frame(
2417 FrameType::Request,
2418 pending.client_channel,
2419 pending.client_epoch + 1,
2420 902,
2421 ),
2422 )
2423 .await
2424 .unwrap();
2425 let stale_error = client_rx.recv().await.unwrap();
2426 assert_eq!(stale_error.header.ty, FrameType::Error);
2427 assert_eq!(stale_error.header.channel, pending.client_channel);
2428 assert_eq!(stale_error.header.epoch, pending.client_epoch + 1);
2429 assert_eq!(stale_error.header.corr, 902);
2430 let body: ErrorBody = serde_json::from_slice(&stale_error.body).unwrap();
2431 assert_eq!(body.code, "stale_route_epoch");
2432 assert!(module_rx.try_recv().is_err());
2433 let counters = router.counters.snapshot();
2434 assert_eq!(counters["client_frames_dropped_stale_route"], 1);
2435 assert_eq!(counters["module_frames_dropped_no_route"], 0);
2436 }
2437
2438 #[tokio::test]
2439 async fn dropped_module_route_goodbye_increments_counter() {
2440 let (
2441 router,
2442 _forwarding,
2443 client_ctx,
2444 mut client_rx,
2445 _module_ctx,
2446 _module_egress_rx,
2447 mut module_rx,
2448 pending,
2449 ) = dynamic_route_fixture(true);
2450 let _ = client_rx.recv().await;
2451 module_rx.close();
2452
2453 router
2454 .route_for_connection(
2455 &client_ctx,
2456 route_frame(
2457 FrameType::Goodbye,
2458 pending.client_channel,
2459 pending.client_epoch,
2460 999,
2461 ),
2462 )
2463 .await
2464 .unwrap();
2465
2466 let counters = router.counters.snapshot();
2467 assert_eq!(counters["goodbye_relay_module_dropped"], 1);
2468 assert_eq!(
2469 counters["goodbye_relay_module_dropped_by_module"],
2470 serde_json::json!({ "epoch-router": 1 })
2471 );
2472 assert_eq!(counters["route_released_epoch_fenced"], 1);
2473 }
2474
2475 #[tokio::test]
2476 async fn module_request_on_stale_epoch_receives_stale_route_epoch() {
2477 let (
2478 router,
2479 _forwarding,
2480 _client_ctx,
2481 _client_rx,
2482 module_ctx,
2483 mut module_egress_rx,
2484 mut module_rx,
2485 pending,
2486 ) = dynamic_route_fixture(true);
2487
2488 router
2489 .route_for_connection(
2490 &module_ctx,
2491 route_frame(
2492 FrameType::Request,
2493 pending.module_channel,
2494 pending.module_epoch + 1,
2495 1_000,
2496 ),
2497 )
2498 .await
2499 .unwrap();
2500
2501 let error = module_egress_rx.try_recv().unwrap();
2502 assert_eq!(error.header.ty, FrameType::Error);
2503 assert_eq!(error.header.channel, pending.module_channel);
2504 assert_eq!(error.header.epoch, pending.module_epoch + 1);
2505 assert_eq!(error.header.corr, 1_000);
2506 let body: ErrorBody = serde_json::from_slice(&error.body).unwrap();
2507 assert_eq!(body.code, "stale_route_epoch");
2508 assert!(module_rx.try_recv().is_err());
2509 let counters = router.counters.snapshot();
2510 assert_eq!(counters["module_requests_dropped_stale_route"], 1);
2511 assert_eq!(counters["module_frames_dropped_no_route"], 0);
2512 }
2513
2514 #[tokio::test]
2515 async fn module_request_on_reserved_or_absent_route_receives_unknown_channel() {
2516 let (
2517 reserved_router,
2518 _forwarding,
2519 _client_ctx,
2520 _client_rx,
2521 reserved_module_ctx,
2522 mut reserved_module_egress_rx,
2523 _module_rx,
2524 reserved,
2525 ) = dynamic_route_fixture(false);
2526 reserved_router
2527 .route_for_connection(
2528 &reserved_module_ctx,
2529 route_frame(
2530 FrameType::Request,
2531 reserved.module_channel,
2532 reserved.module_epoch,
2533 1_001,
2534 ),
2535 )
2536 .await
2537 .unwrap();
2538 let reserved_error = reserved_module_egress_rx.try_recv().unwrap();
2539 let reserved_body: ErrorBody = serde_json::from_slice(&reserved_error.body).unwrap();
2540 assert_eq!(reserved_error.header.ty, FrameType::Error);
2541 assert_eq!(reserved_error.header.channel, reserved.module_channel);
2542 assert_eq!(reserved_error.header.epoch, reserved.module_epoch);
2543 assert_eq!(reserved_error.header.corr, 1_001);
2544 assert_eq!(reserved_body.code, "unknown_channel");
2545 assert_eq!(
2546 reserved_router.counters.snapshot()["module_requests_dropped_stale_route"],
2547 1
2548 );
2549
2550 let (
2551 absent_router,
2552 _forwarding,
2553 _client_ctx,
2554 _client_rx,
2555 absent_module_ctx,
2556 mut absent_module_egress_rx,
2557 _module_rx,
2558 absent,
2559 ) = dynamic_route_fixture(false);
2560 absent_router
2561 .route_for_connection(
2562 &absent_module_ctx,
2563 route_frame(
2564 FrameType::Request,
2565 absent.module_channel + 1,
2566 absent.module_epoch,
2567 1_002,
2568 ),
2569 )
2570 .await
2571 .unwrap();
2572 let absent_error = absent_module_egress_rx.try_recv().unwrap();
2573 let absent_body: ErrorBody = serde_json::from_slice(&absent_error.body).unwrap();
2574 assert_eq!(absent_error.header.ty, FrameType::Error);
2575 assert_eq!(absent_error.header.channel, absent.module_channel + 1);
2576 assert_eq!(absent_error.header.epoch, absent.module_epoch);
2577 assert_eq!(absent_error.header.corr, 1_002);
2578 assert_eq!(absent_body.code, "unknown_channel");
2579 assert_eq!(
2580 absent_router.counters.snapshot()["module_requests_dropped_stale_route"],
2581 1
2582 );
2583 }
2584
2585 #[tokio::test]
2586 async fn non_request_module_frame_on_dead_route_is_counted_without_error() {
2587 let (
2588 router,
2589 forwarding,
2590 client_ctx,
2591 mut client_rx,
2592 module_ctx,
2593 mut module_egress_rx,
2594 mut module_rx,
2595 pending,
2596 ) = dynamic_route_fixture(true);
2597 let (other_module_tx, _other_module_rx) = mpsc::channel(8);
2598 forwarding
2599 .register_module_connection(
2600 ConnectionId::new(502),
2601 "other-module".into(),
2602 2,
2603 Concurrency::ModuleManaged,
2604 FrameSink::new(other_module_tx),
2605 )
2606 .unwrap();
2607 let _ = client_rx.recv().await.unwrap();
2608
2609 router
2610 .route_for_connection(
2611 &client_ctx,
2612 route_frame(
2613 FrameType::Goodbye,
2614 pending.client_channel,
2615 pending.client_epoch,
2616 1_003,
2617 ),
2618 )
2619 .await
2620 .unwrap();
2621 let _ = module_rx.recv().await.unwrap();
2622
2623 router
2624 .route_for_connection(
2625 &module_ctx,
2626 route_frame(
2627 FrameType::StreamData,
2628 pending.module_channel,
2629 pending.module_epoch,
2630 1_004,
2631 ),
2632 )
2633 .await
2634 .unwrap();
2635
2636 let reply = module_egress_rx.try_recv().unwrap();
2639 assert_eq!(reply.header.ty, FrameType::Goodbye);
2640 assert!(module_egress_rx.try_recv().is_err());
2641 let counters = router.counters.snapshot();
2642 assert_eq!(counters["module_frames_dropped_no_route"], 1);
2643 assert_eq!(
2644 counters["module_frames_dropped_no_route_by_module"],
2645 serde_json::json!({ "epoch-router": 1 })
2646 );
2647 assert_eq!(counters["module_requests_dropped_stale_route"], 0);
2648 }
2649
2650 fn drain_now(rx: &mut mpsc::Receiver<OutboundFrame>) {
2653 while rx.try_recv().is_ok() {}
2654 }
2655
2656 #[tokio::test]
2662 async fn route_goodbye_refused_by_stalled_module_is_delivered_when_it_resumes_reading() {
2663 const BUDGET: usize = 4_096;
2664 let forwarding = Arc::new(ForwardingTable::default());
2665 let control = Arc::new(ControlHandler::with_forwarding(
2666 Arc::new(Registry::default()),
2667 Arc::clone(&forwarding),
2668 ));
2669 let router = Router::with_control_handler(control);
2670 let module_connection = ConnectionId::new(80);
2671 let client_connection = ConnectionId::new(81);
2672 let (module_tx, mut module_rx) = mpsc::channel(64);
2673 let module_sink = FrameSink::with_byte_budget(module_tx, BUDGET);
2674 forwarding
2675 .register_module_connection(
2676 module_connection,
2677 "stalled-provider".into(),
2678 2,
2679 Concurrency::ModuleManaged,
2680 module_sink.clone(),
2681 )
2682 .unwrap();
2683 let (client_tx, mut client_rx) = mpsc::channel(8);
2684 let client_sink = FrameSink::new(client_tx);
2685 let pending = forwarding
2686 .begin_route_bind_relay_for_test(
2687 client_connection,
2688 client_sink.clone(),
2689 1_100,
2690 "stalled-provider",
2691 )
2692 .unwrap();
2693 forwarding
2694 .complete_pending_relay(
2695 module_connection,
2696 pending.corr,
2697 RouteBindRelayOutcome::Accepted,
2698 )
2699 .unwrap();
2700 let _ = client_rx.recv().await.unwrap();
2701 drain_now(&mut module_rx);
2702
2703 module_sink
2706 .try_send(stream_frame(9, 1, 0, vec![b'f'; BUDGET]))
2707 .unwrap();
2708 assert!(module_sink
2709 .try_send(stream_frame(9, 1, 1, Vec::new()))
2710 .is_err());
2711
2712 let client_ctx = RouteCtx {
2713 connection_id: client_connection,
2714 egress: client_sink,
2715 };
2716 router
2717 .route_for_connection(
2718 &client_ctx,
2719 route_frame(
2720 FrameType::Goodbye,
2721 pending.client_channel,
2722 pending.client_epoch,
2723 1_101,
2724 ),
2725 )
2726 .await
2727 .unwrap();
2728 tokio::task::yield_now().await;
2729 assert_eq!(
2730 router.counters.snapshot()["goodbye_relay_module_dropped"],
2731 0,
2732 "a GOODBYE refused by a momentarily full module queue must not be dropped"
2733 );
2734
2735 let filler = module_rx.recv().await.unwrap();
2737 assert_eq!(filler.header.ty, FrameType::StreamData);
2738 drop(filler);
2739 let goodbye = tokio::time::timeout(Duration::from_secs(2), module_rx.recv())
2740 .await
2741 .expect("the refused GOODBYE must be delivered once the module frees room")
2742 .unwrap();
2743 assert_eq!(goodbye.header.ty, FrameType::Goodbye);
2744 assert_eq!(goodbye.header.channel, pending.module_channel);
2745 assert_eq!(goodbye.header.epoch, pending.module_epoch);
2746 assert_eq!(
2747 router.counters.snapshot()["goodbye_relay_module_dropped"],
2748 0
2749 );
2750 }
2751
2752 #[tokio::test]
2758 async fn module_frame_on_route_the_daemon_does_not_hold_is_answered_with_goodbye() {
2759 let (
2760 router,
2761 _forwarding,
2762 client_ctx,
2763 mut client_rx,
2764 module_ctx,
2765 mut module_egress_rx,
2766 mut module_rx,
2767 pending,
2768 ) = dynamic_route_fixture(true);
2769 let _ = client_rx.recv().await.unwrap();
2770 router
2771 .route_for_connection(
2772 &client_ctx,
2773 route_frame(
2774 FrameType::Goodbye,
2775 pending.client_channel,
2776 pending.client_epoch,
2777 1_200,
2778 ),
2779 )
2780 .await
2781 .unwrap();
2782 drain_now(&mut module_rx);
2783
2784 let cases = [
2785 (pending.module_channel, pending.module_epoch),
2787 (pending.module_channel + 1, 1),
2789 ];
2790 for (channel, epoch) in cases {
2791 router
2792 .route_for_connection(
2793 &module_ctx,
2794 route_frame(FrameType::StreamData, channel, epoch, 1_201),
2795 )
2796 .await
2797 .unwrap();
2798 let reply = module_egress_rx
2799 .try_recv()
2800 .expect("a frame on a route the daemon does not hold is answered");
2801 assert_eq!(reply.header.ty, FrameType::Goodbye);
2802 assert_eq!(reply.header.channel, channel);
2803 assert_eq!(reply.header.epoch, epoch);
2804 assert_eq!(reply.header.corr, 0);
2805 assert!(module_egress_rx.try_recv().is_err());
2806 }
2807
2808 router
2811 .route_for_connection(
2812 &module_ctx,
2813 route_frame(FrameType::Goodbye, pending.module_channel + 2, 1, 0),
2814 )
2815 .await
2816 .unwrap();
2817 assert!(module_egress_rx.try_recv().is_err());
2818
2819 let counters = router.counters.snapshot();
2820 assert_eq!(counters["module_frames_dropped_no_route"], 3);
2821 assert_eq!(counters["module_frames_dropped_released_route"], 1);
2822 assert_eq!(
2823 counters["module_frames_dropped_released_route_by_module"],
2824 serde_json::json!({ "epoch-router": 1 })
2825 );
2826 assert_eq!(counters["module_orphan_route_goodbyes_sent"], 2);
2827 }
2828
2829 #[tokio::test(start_paused = true)]
2833 async fn burst_of_orphan_module_frames_is_answered_once_per_interval() {
2834 let (
2835 router,
2836 _forwarding,
2837 client_ctx,
2838 mut client_rx,
2839 module_ctx,
2840 mut module_egress_rx,
2841 mut module_rx,
2842 pending,
2843 ) = dynamic_route_fixture(true);
2844 let _ = client_rx.recv().await.unwrap();
2845 router
2846 .route_for_connection(
2847 &client_ctx,
2848 route_frame(
2849 FrameType::Goodbye,
2850 pending.client_channel,
2851 pending.client_epoch,
2852 1_300,
2853 ),
2854 )
2855 .await
2856 .unwrap();
2857 drain_now(&mut module_rx);
2858
2859 let send_burst = |corr: u64| {
2860 route_frame(
2861 FrameType::StreamData,
2862 pending.module_channel,
2863 pending.module_epoch,
2864 corr,
2865 )
2866 };
2867 for corr in 0..20 {
2868 router
2869 .route_for_connection(&module_ctx, send_burst(corr))
2870 .await
2871 .unwrap();
2872 }
2873 let mut replies = 0;
2874 while let Ok(reply) = module_egress_rx.try_recv() {
2875 assert_eq!(reply.header.ty, FrameType::Goodbye);
2876 replies += 1;
2877 }
2878 assert_eq!(replies, 1, "a burst within the interval gets one GOODBYE");
2879
2880 tokio::time::advance(ORPHAN_ROUTE_GOODBYE_INTERVAL).await;
2881 router
2882 .route_for_connection(&module_ctx, send_burst(20))
2883 .await
2884 .unwrap();
2885 let retry = module_egress_rx
2886 .try_recv()
2887 .expect("the first orphan frame after the interval is answered again");
2888 assert_eq!(retry.header.channel, pending.module_channel);
2889 assert_eq!(retry.header.epoch, pending.module_epoch);
2890
2891 let counters = router.counters.snapshot();
2892 assert_eq!(counters["module_frames_dropped_no_route"], 21);
2893 assert_eq!(counters["module_orphan_route_goodbyes_sent"], 2);
2894 }
2895
2896 #[test]
2898 fn orphan_goodbye_rate_limit_state_is_released_with_the_connection() {
2899 let router = Router::with_default_self_handler();
2900 let connection = router.begin_connection();
2901 let id = connection.id();
2902 assert!(router.orphan_goodbyes.claim(id, 7));
2903 assert!(!router.orphan_goodbyes.claim(id, 7));
2904 assert!(router.orphan_goodbyes.claim(id, 8));
2905 drop(connection);
2906 assert!(router
2907 .orphan_goodbyes
2908 .last_sent
2909 .lock()
2910 .unwrap()
2911 .get(&id)
2912 .is_none());
2913 }
2914
2915 #[test]
2916 fn channel_zero_cannot_be_registered_as_backend() {
2917 let mut router = Router::with_default_self_handler();
2918
2919 let err = router.register_backend(0, EchoBackend).unwrap_err();
2920
2921 assert_eq!(err, RouterError::ReservedChannelZero);
2922 }
2923}