Skip to main content

subc_daemon/
forwarding.rs

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
28/// Default per-channel request-credit window for modules that schedule internally.
29const DEFAULT_MODULE_MANAGED_WINDOW: usize = 32;
30
31/// High per-channel cap for stateless modules; this is an OOM guard, not scheduling policy.
32const STATELESS_PARALLEL_WINDOW: usize = 1024;
33
34/// A stopped probe cycle cannot retain its last unanswered correlation forever.
35/// Active endpoints replace the tombstone on their next serial health probe;
36/// this backstop covers endpoints that stop probing altogether.
37const HEALTH_PROBE_TOMBSTONE_TTL: Duration = Duration::from_secs(5 * 60);
38
39/// Module connection identity used in forwarding keys.
40///
41/// The generation is bumped every time a module connection is registered so a future restart cannot
42/// accidentally reuse a stale `(connection_id, route_channel)` binding.
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
44pub struct ModuleEndpointId {
45    pub connection_id: ConnectionId,
46    pub generation: u64,
47}
48
49/// Client-local route key. A route channel is unique only within one client connection.
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
51pub(crate) struct ClientRouteKey {
52    pub connection_id: ConnectionId,
53    pub channel: u16,
54}
55
56/// Module-local route key. A route channel is unique only within one live module endpoint.
57#[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    /// The scope the route was admitted under, with the tag its bind was
81    /// stamped at. A scope change finds the routes to close by it.
82    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/// Which kind of peer a route GOODBYE is being delivered to. This decides what
100/// happens when the GOODBYE cannot be enqueued (egress full/closed):
101/// - `Client`: escalate to closing that client connection (a socket close is a
102///   stronger teardown signal, and a full client egress means it is the slow
103///   client we would drop anyway).
104/// - `Module`: never close; deliver late instead (see
105///   [`send_module_route_goodbye`]), and drop only if the module still has no
106///   room after [`LATE_MODULE_GOODBYE_DEADLINE`]. A client-disconnect notifies the
107///   SHARED module that one client's route is gone; closing the module on its
108///   egress backpressure would tear down every co-tenant client (the exact
109///   cross-tenant blast radius this never-close rule exists to prevent — observed when a
110///   flooding dead client filled BOTH its own and the module's egress, so its
111///   route-gone GOODBYE to the module failed and closed the shared connection).
112///   subc has already removed the route from its forwarding state and drops
113///   stale module frames for the released channel (see router.rs), so subc's
114///   own routing is correct. The residual: under SUSTAINED module-egress
115///   backpressure a module-targeted route-gone notification can be lost, which a
116///   module using it for client-refcounting (e.g. AFT's session accounting)
117///   would miss. This is INTENTIONALLY ACCEPTED, not a gap. A consuming module
118///   must bound stale bindings with its own idle-activity reaper (last-touched
119///   TTL) independent of route-gone signals — AFT does exactly this, so a lost
120///   GOODBYE degrades to "the
121///   binding stays warm until its idle TTL" (bounded wasted resources), never an
122///   unbounded leak; disk-durable replay is unaffected. A dedicated reliable
123///   module control lane was evaluated and deliberately NOT
124///   built: it would add starvation-avoidance machinery to the thin core for a
125///   bounded warm-resource window that is not a correctness issue. never-close
126///   is the invariant that matters here.
127#[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    /// The module on the other end of the route: for a module-targeted relay
142    /// the receiving module (attributing a dropped relay), for a client target
143    /// the module whose route is going away (naming it when an undeliverable
144    /// relay closes the client).
145    pub module_id: Option<String>,
146}
147
148/// The frame a client connection's egress queue refused, for the diagnosis
149/// logged when that refusal closes the connection.
150#[derive(Debug, Clone, Copy)]
151pub(crate) struct UndeliveredFrame<'a> {
152    /// The module on the other end of the route the frame belonged to, when known.
153    pub module_id: Option<&'a str>,
154    /// The refusing connection's sink, read for what its queue held.
155    pub sink: &'a FrameSink,
156}
157
158/// How a route's principal appears in logs: `direct`, or `reserved:<module>`.
159fn 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
167/// The distinct principals of the routes bound on one client connection,
168/// sorted and comma-joined, or `none` when it has no bound route.
169fn 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    /// True only when an undeliverable GOODBYE should escalate to closing the
185    /// target connection. Never escalate for module recipients.
186    pub(crate) fn close_on_delivery_failure(&self) -> bool {
187        matches!(self.kind, GoodbyeTargetKind::Client)
188    }
189}
190
191/// How long a route GOODBYE that a module's egress queue refused keeps waiting
192/// for room before it is given up and counted as dropped.
193///
194/// A module that stops reading for a moment (a GC pause, a long synchronous
195/// handler, a reconnect herd it is still working through) should still learn
196/// that the route is gone, because a module that never hears the GOODBYE keeps
197/// the route and keeps sending on it for as long as its connection lives. The
198/// value is the default module drain timeout: the daemon already treats that as
199/// the longest it is reasonable to wait on a module that is busy but alive, and
200/// a module that cannot free 21 bytes of egress in that time is wedged rather
201/// than slow. Resolving each module's own configured drain timeout here would
202/// need a registry lookup on a path that holds only the module's sink.
203pub(crate) const LATE_MODULE_GOODBYE_DEADLINE: Duration = crate::supervise::DEFAULT_DRAIN_TIMEOUT;
204
205/// Send a route GOODBYE to a module without ever closing the module's shared
206/// connection, which would sever every other route it serves.
207///
208/// The frame is enqueued at once when the module's egress queue has room. When
209/// it does not, a detached task waits for room (bounded by
210/// [`LATE_MODULE_GOODBYE_DEADLINE`]) and sends it then. Arriving late is safe:
211/// the daemon has already released the route, and module SDKs tear a route
212/// down only when both the channel AND the epoch match the route installed on
213/// that channel, so a late GOODBYE for a channel since reused at a newer epoch
214/// is ignored. The GOODBYE is counted in `goodbye_relay_module_dropped` only
215/// when the connection is closed or the deadline passes.
216pub(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    // Waiting is pointless once the module's writer is gone, and impossible
229    // outside a Tokio runtime (cleanup that runs from a destructor at shutdown).
230    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/// One route currently served by a module endpoint.
287///
288/// `goodbye_target` is deliberately retained alongside the census projection so
289/// a draining caller can address exactly the same route set that this read
290/// reports, without a second forwarding-table pass.
291#[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    /// WHY the endpoint is draining, when it is. Carried per-route so the
298    /// census can answer "closing because of what" without a second lookup;
299    /// `None` exactly when `draining` is false (one source: the drain map).
300    pub drain_reason: Option<RouteCloseReason>,
301}
302
303/// One live route a scope change closed: both ends to send GOODBYE to, and the
304/// reason for the client's `route.closed`.
305#[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/// One registered module connection, as [`ForwardingTable::module_connections`]
336/// reports it: enough to send it a frame and then close it.
337#[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/// What [`ForwardingTable::cutover_candidate`] did.
354#[derive(Debug, Clone, Copy, PartialEq, Eq)]
355pub(crate) struct ForwardingCutover {
356    /// The endpoint now in the active slot (the former candidate).
357    pub promoted: ModuleEndpointId,
358    /// The endpoint demoted out of the active slot, to be drained by endpoint.
359    /// `None` when the incumbent's connection was already gone.
360    pub incumbent: Option<ModuleEndpointId>,
361}
362
363/// What tearing down one connection released.
364#[derive(Debug)]
365pub(crate) struct ConnectionCleanup {
366    /// Routes whose other end must be sent GOODBYE.
367    pub released: Vec<GoodbyeTarget>,
368    /// Pending route.bind relays to the closed module that were aborted: opens
369    /// that had not been bound yet. Zero when a client connection closed.
370    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    /// The scope admitted at `route.open`, with the tag captured then. Commit
442    /// compares it with the published tag.
443    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    /// The ACTIVE slot: the one endpoint per module id that routing resolves.
478    /// Every by-id lookup (relay reservation, drain-by-id, liveness, census,
479    /// live roots, module-control RPCs) reads this map and nothing else.
480    modules_by_id: HashMap<String, ModuleConnection>,
481    /// The CANDIDATE slot of a blue/green swap: a second process registered
482    /// under an id that already has an active endpoint. It has a full endpoint
483    /// identity (so its own connection can be looked up and torn down) but no
484    /// by-id lookup sees it, so nothing is routed to it until `cutover_candidate`
485    /// promotes it.
486    candidates_by_id: HashMap<String, ModuleConnection>,
487    /// Former active endpoints demoted by `cutover_candidate`, until their
488    /// connection is removed. Membership is what distinguishes an endpoint that
489    /// was SUPERSEDED by a promotion (still alive, still carrying bound routes,
490    /// about to be drained) from one that is merely STALE (replaced some other
491    /// way, which is treated as a fault on the acking connection).
492    superseded_endpoints: HashMap<ModuleEndpointId, ModuleConnection>,
493    endpoint_by_connection: HashMap<ConnectionId, ModuleEndpointId>,
494    module_id_by_endpoint: HashMap<ModuleEndpointId, String>,
495    /// Endpoints mid-drain, keyed to the reason the drain was begun with. The
496    /// value serves the census ("draining because restart"); membership alone
497    /// still answers every admission-gate check.
498    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    /// Each scope's current `(scope_epoch, version)`, keyed `(owner, ref)`,
513    /// published by a scope sync in the same step that changes the record. A
514    /// bind commit reads it here so it never takes the scope table's lock
515    /// while holding this one. Absent means the scope is not live.
516    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/// Dynamic forwarding state shared by the control plane and data-plane router.
546#[derive(Debug, Default)]
547pub struct ForwardingTable {
548    inner: Arc<RwLock<ForwardingInner>>,
549    close_registry: Mutex<HashMap<ConnectionId, oneshot::Sender<CloseReason>>>,
550    counters: DaemonCounters,
551    /// Per-target-module bind-relay breaker state. It lives here, beside the
552    /// module connections it describes, because this is where a module
553    /// connection's identity is established and therefore where a stale
554    /// verdict has to be discarded.
555    route_bind_breakers: RouteBindBreakers,
556    /// Current route.bind relays keyed by target module. Admission is shared
557    /// across every client connection that points at the same endpoint.
558    route_bind_concurrency: RouteBindConcurrency,
559    /// Start and end of each module's `route.open` outage. Held here for the
560    /// same reason as the breakers: every handler built over this table must
561    /// see one outage per module, or each would log its own opening line.
562    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    /// Keep the binding protected through confirmation admission. Route release
593    /// takes the write lock and then the same confirm-state lock, so it cannot
594    /// miss an admitted request or leave it waiting on a vanished route.
595    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    /// Close every established connection, modules and clients alike, as the
655    /// last step of an announced daemon shutdown. By the time this runs each
656    /// registered module has already been sent a module GOODBYE (see
657    /// `Supervisor::end_children_for_daemon_shutdown`), so the EOF that follows
658    /// is a planned stop; EOF with no GOODBYE before it stays reserved for a
659    /// daemon that went away unannounced. Closing while the daemon is still
660    /// alive lets it wait for the modules' own teardowns before it ends
661    /// whatever is left.
662    #[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    /// Every registered module connection, whichever slot it occupies (active,
673    /// swap candidate, or superseded incumbent), once each. Daemon shutdown
674    /// uses this to tell each of them it is a planned stop before closing it.
675    #[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    /// Ask a registered connection to close. Returns true only for the request
701    /// that actually reached it; later requests for the same connection, and
702    /// requests for a connection that is not registered, return false.
703    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    /// Register a module connection and queue its HELLO_ACK as the first frame
746    /// on its sink, in the same write-lock critical section that makes the
747    /// endpoint visible.
748    ///
749    /// A module reads HELLO_ACK as the first frame after its HELLO and exits on
750    /// anything else. Once the endpoint is in `modules_by_id`, a `route.open`
751    /// on any other connection can queue a `route.bind` request onto this sink,
752    /// so the ack must already be queued by then. Every lookup that can queue a
753    /// frame for the module takes this lock, so nothing can get in ahead of it.
754    ///
755    /// If the ack cannot be queued (the sink is closed or full), the
756    /// registration fails with nothing inserted.
757    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        // Every refusal check has passed and nothing has been mutated yet, so
791        // a failed enqueue leaves the table exactly as it was.
792        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        // A new module connection has arrived under this id, so anything the
817        // bind-relay breaker learned was learned about a process that is no
818        // longer the one behind this name. See
819        // `RouteBindBreakers::reset_for_new_module_connection`.
820        //
821        // Keyed on ARRIVAL rather than on teardown deliberately: a connection
822        // going away is not evidence about anything, and a module whose
823        // connection drops without coming back should keep its verdict until
824        // something actually registers in its place. This also covers the
825        // unclean replacements -- a module killed mid-bind, or one whose
826        // connection was closing -- because registration is the single path by
827        // which any module connection becomes usable.
828        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    /// Register a blue/green swap candidate for `module_id` into the candidate
842    /// slot, alongside whatever endpoint is active for the id.
843    ///
844    /// Unlike [`Self::register_module_connection`], this never touches the
845    /// active slot and never resets the bind-relay breaker: the incumbent is
846    /// still the process serving the id, so what the breaker learned about it
847    /// is still true. The breaker is reset at cutover instead, when the process
848    /// behind the name actually changes. A second candidate for the same id is
849    /// refused.
850    ///
851    /// Production registers candidates through
852    /// [`Self::register_candidate_module_connection_acked`]; this ack-less form
853    /// is for tests that build forwarding state directly.
854    #[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    /// The swap-candidate counterpart of
874    /// [`Self::register_module_connection_acked`]: the HELLO_ACK is queued on
875    /// the candidate's sink before its endpoint is inserted. No by-id lookup
876    /// sees a candidate, but its own connection's endpoint does become
877    /// resolvable here, and at cutover it becomes routable; the ack has to be
878    /// ahead of anything either can queue.
879    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        // As in `register_module_connection_inner`: after every refusal check,
916        // before any mutation.
917        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    /// Promote `module_id`'s candidate to the active slot, in ONE forwarding
943    /// write-lock critical section.
944    ///
945    /// That single section is the linearization point of a swap. Every relay
946    /// reservation resolves the active slot under the same lock, so each one
947    /// lands wholly before cutover (on the incumbent) or wholly after it (on the
948    /// promoted candidate), never on a mix. A relay reserved on the incumbent
949    /// before cutover can no longer commit: `commit_route_locked` requires the
950    /// reservation's endpoint to be the active one, and an ack for it arriving
951    /// later is answered by the superseded-endpoint arm of
952    /// `complete_pending_relay` instead of ending the incumbent's connection.
953    ///
954    /// Both endpoints keep their identities: nothing is re-keyed, so the
955    /// incumbent's bound routes and its pending correlation keys stay exactly
956    /// where they are until it is drained with [`Self::begin_endpoint_drain`],
957    /// using the incumbent endpoint this returns. Returns `Ok(None)` when there
958    /// is no candidate for the id.
959    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        // A different process now answers for this id; see the matching reset
982        // in `register_module_connection` for why a verdict about the old one
983        // must not carry over.
984        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        // Reserve egress capacity before taking the forwarding lock. The permit is
1014        // held until the bind reaches one terminal state, so an accepted bind can
1015        // publish its table entry and RouteOpen response in one critical section.
1016        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    /// A health probe addressed to one endpoint, in whichever slot it is.
1107    ///
1108    /// Every other control RPC resolves the id's ACTIVE endpoint, which is how
1109    /// a swap candidate (not active until cutover) and a superseded incumbent
1110    /// (not active after it) are unreachable by them. A swap needs to probe
1111    /// exactly those two: the candidate before promoting it, and the incumbent
1112    /// for its busy gauges while it drains.
1113    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        // A draining incumbent is exactly what this probes for busy gauges, so
1130        // draining is allowed, as it is for the by-id drain probe.
1131        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            // Recurring health probes are serial per endpoint. Once the next one
1193            // starts, an older answer can no longer improve the current snapshot,
1194            // so retaining more than the newest unanswered probe has no value.
1195            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        // HELLO and route.open can run on different tasks. Decide the role under
1249        // the same lock that installs endpoints and reserves client channels so
1250        // neither task can turn one connection into both kinds of owner.
1251        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            // Two shapes of "the client is not there to receive this route" that
1555            // must resolve identically: its egress is already closed, or it is
1556            // marked closing (its connection loop has been asked to end, but has
1557            // not drained yet, so the sink is still open).
1558            //
1559            // Only the first used to be caught here. The second fell through to
1560            // `commit_route_locked`, which refuses a closing client with
1561            // `ConnectionClosing` -- and this function is called from the MODULE
1562            // connection's frame handler, so that refusal ended the module's
1563            // connection instead of this one client's route. One dying client's
1564            // route.open then took down a connection carrying every other
1565            // client's routes to that module.
1566            //
1567            // The remedy for both is the same, which is why they share an arm:
1568            // give back the reserved handle pair, tell the waiting route.open the
1569            // route is gone, and report the module-side channel so the caller can
1570            // send a channel-scoped GOODBYE for the binding the module just
1571            // created. Nothing here touches the module connection.
1572            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            // The endpoint that acked this bind is no longer accepting new
1601            // routes, for one of two reasons. Either it was the module's active
1602            // endpoint when this relay was reserved and a blue/green swap has
1603            // since promoted a candidate (the replacement process registered
1604            // beside it) into the active slot, superseding it; or a daemon
1605            // drain has gated it, and the ack arrived before the ordered module
1606            // drain settled pending relays. The module did nothing wrong: it
1607            // bound the route it was asked to bind, and the endpoint still
1608            // carries every other client's routes until it is drained. So the
1609            // ack is settled here, while the pending entry still holds the
1610            // client's sender and the reserved handle pair: release the pair,
1611            // refuse the waiting route.open with `module_reloading`, which
1612            // callers retry (after a swap the retry reserves on the promoted
1613            // endpoint), and hand back the module-side channel so the caller
1614            // sends one channel-scoped GOODBYE for the binding the module just
1615            // created. Returning an error instead would end the module
1616            // connection all those other routes share; nothing here touches it.
1617            //
1618            // An endpoint that stopped being active any other way is stale
1619            // (replaced, but not by a swap and not draining). It still falls
1620            // through to `commit_route_locked`, which refuses it with
1621            // `StaleModuleEndpoint` exactly as before swaps existed.
1622            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                // Not gated on `relay_enqueued`: an ack proves the module
1632                // received the bind, whether or not the enqueue mark was set.
1633                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            // The scope the open was admitted under ended or changed after
1651            // admission. The module bound a route under a stamp that is no
1652            // longer true, so it must not become routable. Settled here like
1653            // the superseded arm above, never as an error from
1654            // `commit_route_locked`, which would close the module's whole
1655            // connection and every other client's routes to it: release the
1656            // pair, answer the waiting route.open by name, and hand back the
1657            // module-side channel for one channel-scoped GOODBYE.
1658            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    /// Reserved-but-uncommitted route handles, counted on both index sides, so
1794    /// a test can assert a reservation pair was actually given back.
1795    #[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    /// Looks up the module registered on a data-plane connection so route-drop
1813    /// diagnostics name the emitter instead of only reporting a daemon total.
1814    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    /// Whether the daemon ever allocated `(channel, epoch)` on this module
1827    /// connection. Epochs on a module channel are handed out as 1, 2, 3, ...
1828    /// and the last one handed out is remembered until the connection ends, so
1829    /// every epoch from 1 up to that one was a real route at some point. Used
1830    /// only on the drop path, to tell a module still sending on a route the
1831    /// daemon released from one sending on a route that never existed.
1832    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    /// How many distinct client connections hold at least one committed route,
2031    /// alongside the largest number of routes any single connection holds.
2032    ///
2033    /// `connected_clients` alone cannot distinguish many clients with a route
2034    /// each from one client accumulating hundreds, and those have opposite
2035    /// causes. Reading it required an out-of-band `lsof` during a live
2036    /// investigation, and the count of connections was mistaken for a count of
2037    /// client processes — which sent two of us after cleanup paths that were
2038    /// working correctly.
2039    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    /// Whether daemon-wide shutdown has begun draining all providers.
2058    pub(crate) fn is_daemon_draining(&self) -> Result<bool, ForwardingError> {
2059        Ok(self.read_inner()?.daemon_draining)
2060    }
2061
2062    /// Gate every provider atomically, including registrations racing shutdown.
2063    /// No supervisor lock is held while taking the forwarding lock.
2064    #[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        // Swap candidates and superseded incumbents are not routable, but they
2079        // are live endpoints that could still be sent module-control work, so
2080        // the daemon-wide gate covers them too. The returned ids are unchanged:
2081        // they name modules for the supervisor to drain, one per id.
2082        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    /// Begin draining whatever endpoint is ACTIVE for `module_id`.
2097    ///
2098    /// After a swap's cutover the active endpoint is the promoted candidate, so
2099    /// this must not be used to drain the incumbent; use
2100    /// [`Self::begin_endpoint_drain`] with the endpoint `cutover_candidate`
2101    /// returned.
2102    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    /// Begin draining one specific endpoint, whichever slot it is in.
2117    ///
2118    /// This is how a swap drains its incumbent after cutover: the incumbent's
2119    /// endpoint is captured by `cutover_candidate`, and resolving it by module id
2120    /// instead would find the promoted candidate and leave neither process
2121    /// routable. Returns `Ok(None)` when the endpoint is no longer registered.
2122    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
2142/// Mark `module`'s endpoint draining and settle everything still pending on it.
2143/// Shared by the by-id and by-endpoint drain entry points, which differ only in
2144/// how they find the endpoint.
2145fn 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/// What a drain that timed out was still waiting on, for the log line that
2222/// reports the timeout. Per-request ages are not tracked.
2223#[derive(Debug, Default, PartialEq, Eq)]
2224pub(crate) struct DrainHoldouts {
2225    /// Requests still counted against the drain (subscriptions flagged as
2226    /// such are excluded from the drain and not counted here).
2227    pub(crate) requests: usize,
2228    /// Routes holding at least one of those requests.
2229    pub(crate) routes: usize,
2230    /// Every route on the endpoint, for scale.
2231    pub(crate) total_routes: usize,
2232    /// The client connections holding the most requests, largest first, at
2233    /// most three: enough to name the consumer without listing every route.
2234    pub(crate) top_connections: Vec<(u64, usize)>,
2235    /// The held requests themselves, as the module saw them, so its own log
2236    /// can say what each one was: `(module channel, corr)`, ordered by channel
2237    /// then corr, at most [`DRAIN_HELD_REQUESTS_LISTED`]. The daemon never reads
2238    /// request bodies, so it cannot name a request's method; the module can,
2239    /// from the channel and corr. A request is released only when the module
2240    /// sends a terminal frame (Response, Error or StreamEnd) with its corr on
2241    /// its route, so every pair listed here is one the module ended, if at all,
2242    /// without sending that frame.
2243    pub(crate) held: Vec<(u16, u64)>,
2244}
2245
2246/// How many held requests the drain-timeout line lists by channel and corr.
2247pub(crate) const DRAIN_HELD_REQUESTS_LISTED: usize = 32;
2248
2249impl ForwardingTable {
2250    /// Summarise the requests one endpoint's drain is still waiting on.
2251    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    /// Enumerate one endpoint's current routes without contacting the module.
2341    ///
2342    /// The read lock makes this safe while the endpoint drains: the returned
2343    /// `draining` marker describes the same table state that owns the route,
2344    /// rather than inferring liveness from a module that may be stopping.
2345    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    /// Snapshot all live endpoint route sets under one forwarding-table read lock.
2354    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    /// Snapshot only the endpoint currently routable under this module id.
2373    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    /// True if this connection already owns committed or reserved CLIENT routes.
2428    /// A module registers (HELLO) before serving and never opens client routes, so
2429    /// a connection that has client routes must not also become a module endpoint
2430    /// — otherwise one connection holds both client and module state and cleanup
2431    /// only releases one side.
2432    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    /// [`Self::cleanup_connection`], also reporting how many pending
2449    /// route.bind relays to the closed module were aborted. That count is the
2450    /// `abandoned` figure of the `route.closed` push sent for a lost module
2451    /// connection; it is always zero for a client connection.
2452    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        // The closing mark refuses new work for a connection whose teardown is
2467        // still pending. This is the latest point at which lifting it is safe:
2468        // teardown has just removed every per-connection entry above, under
2469        // this same write lock, so no lookup can still find live state for the
2470        // id; and ids come from a monotonic counter that never reissues one,
2471        // so the id can never name a future connection either. Keeping the
2472        // mark past this point only grew the set by one entry per connection
2473        // for the life of the daemon.
2474        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    /// Close a client connection whose egress queue refused a frame for the
2546    /// route `(connection_id, channel)` at `expected_epoch`. The whole
2547    /// connection closes because its queue is shared by every route on it.
2548    ///
2549    /// The first request that actually closes the connection logs one WARN with
2550    /// the diagnosis: which principals the connection's routes belonged to,
2551    /// which module's frame did not fit and on which client channel, and what
2552    /// the queue held at that moment.
2553    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    /// Publish the tags a scope sync changed and close the live routes its
2599    /// drain table selects, in one critical section.
2600    ///
2601    /// The caller holds the scope table's write lock (scope table, then this
2602    /// table, always). Publishing and selecting under one lock is what closes
2603    /// the race with a bind commit: a commit before this sees the old tag and
2604    /// its route is selected here like any other; a commit after it sees the
2605    /// new tag and is refused.
2606    ///
2607    /// Every client route index entry is scanned, which includes routes on a
2608    /// swap's superseded endpoint (the old process of a module being replaced
2609    /// blue/green, which keeps its existing routes until they drain): those
2610    /// stay indexed until drained, so ending a scope reaches them too.
2611    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    /// The published tag of one scope, for tests of the commit re-check.
2691    #[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
2890/// Why a bind admitted under `scope` may no longer commit, read from the
2891/// published tags: `scope_ended` when the scope is gone or at another epoch,
2892/// `scope_changed` (retryable) when only its content moved.
2893fn 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    // Sending through the reserved slot cannot fail, but it reports a receiver
3065    // that closed after reservation and before this locked publication point.
3066    // The frame is stamped and charged to the queued-byte count here, when it
3067    // actually enters the queue, not when the slot was reserved.
3068    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
3101/// The live module connection registered under exactly `endpoint`, whichever
3102/// slot holds it: active, swap candidate, or superseded incumbent.
3103///
3104/// Resolving by endpoint rather than by module id is what keeps an incumbent
3105/// addressable after a swap promoted a candidate over its id. An endpoint that
3106/// is in none of the three (a stale one, replaced without a promotion) resolves
3107/// to nothing, as it always has.
3108fn 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
3145/// Queue a registering module's HELLO_ACK on its sink. Called with the
3146/// forwarding write lock held, before the endpoint is inserted, which is what
3147/// puts the ack ahead of any `route.bind` or control RPC routed to the module.
3148/// `try_send` never waits, so holding the lock across it is safe.
3149fn 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    // Commit the disconnect before its routes are released, so their teardown
3166    // cannot replace module_closed with route_closed.
3167    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    /// The correlation ids of the requests a drain still waits on, one entry
3321    /// per held credit, in ascending order.
3322    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/// Per-channel request-credit accounting shared by the client and module route halves.
3345#[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            // Protocol-conforming modules emit exactly one terminal per request.
3399            // This guard is a best-effort safety net against window growth, not a
3400            // security boundary against malicious peers.
3401            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    /// A connection cannot hold both client routes and a module endpoint, or
3484    /// replace a module identity while that endpoint is still registered.
3485    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    /// A swap candidate is already registered for this module id.
3512    CandidateSlotOccupied {
3513        module_id: String,
3514    },
3515    /// A registering module's HELLO_ACK could not be queued because its
3516    /// outbound queue is closed or full.
3517    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    /// An acked registration whose HELLO_ACK cannot be queued must leave no
3637    /// endpoint behind: a routable module that never got its ack would read a
3638    /// route.bind first and exit. A closed queue and a full one both refuse.
3639    #[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    /// Both acked registration forms put the HELLO_ACK on the module's queue
3715    /// by the time the endpoint can be resolved.
3716    #[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    /// A pending route.open reserves its queue slot with `reserve_owned` and
4121    /// sends its prebuilt response through that slot when the module accepts.
4122    /// A queue already holding many data frames (far more than the old 64-frame
4123    /// queue) must not stop the open from reserving or completing, and the
4124    /// response must be charged to the queue's byte count when it is sent.
4125    #[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    /// The drain-timeout line reports these numbers, so they must count the
4196    /// requests the drain is waiting on and no others: a route with nothing in
4197    /// flight is not a holdout, and a flagged subscription is excluded from the
4198    /// drain and so from the count.
4199    #[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                // The held requests by the module's channel and corr, ascending,
4239                // without the excluded subscription (corr 3).
4240                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        // The distinction this asserts is the one a bare connection count cannot
4592        // make: two connections holding one route each and one connection
4593        // holding two are the same total, and have opposite causes.
4594        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        // One connection, two routes — not two connections with a route each.
4611        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    /// A client can be marked closing while its egress is still open: the daemon
4716    /// asks a connection to close (here through the production path, a module
4717    /// frame that its egress refused) and the connection loop tears down a moment
4718    /// later. A route.bind ack that lands inside that window is answered on the
4719    /// MODULE connection's frame handler, so resolving it must not produce an
4720    /// error -- an error there ends the module connection, and that connection
4721    /// carries every other client's routes to the module.
4722    #[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        // A published route on this client: escalate_client_delivery_failure only
4728        // marks a connection closing for a route it has already published.
4729        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        // A second route.open from the same client, relayed and awaiting its ack.
4740        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        // The window: closing, but the sink is still open.
4746        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        // No route was published to a client that is on its way out, and the
4784        // reserved handle pair went back.
4785        assert!(client_rx.try_recv().is_err());
4786        assert_eq!(forwarding.active_binding_count().unwrap(), 1);
4787
4788        // The module endpoint is untouched: still live, and still able to take a
4789        // route from another client.
4790        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    /// Every connection teardown passes through `cleanup_connection`, and
4876    /// connection ids come from a monotonic counter that never hands an id out
4877    /// twice. If the closing mark survives teardown, the set grows by one entry
4878    /// per connection for the life of the daemon -- the self-watchdog alone
4879    /// reconnects once a minute.
4880    #[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    /// The closing mark exists to refuse new work for a connection that is on
4911    /// its way out but whose teardown has not run yet: the daemon asks the
4912    /// connection loop to end, and only when the loop reacts does
4913    /// `cleanup_connection` strip the connection's state. Inside that window an
4914    /// operation for the dying connection must still be refused; only after
4915    /// cleanup completes may the mark go.
4916    #[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        // A published route: escalate_client_delivery_failure only marks a
4922        // connection closing for a route it has already published.
4923        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        // Mark the connection closing through the production path without
4934        // running teardown, pinning the window open.
4935        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        // A late route.open for the closing client is refused ...
4953        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        // ... and so is a late attempt to register the connection as a module.
4959        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        // Teardown is the point that lifts the mark: it has just removed every
4973        // per-connection entry under the same lock, so the gate has nothing
4974        // left to protect for this id.
4975        forwarding.cleanup_connection(client).unwrap();
4976        assert_eq!(forwarding.closing_connection_count().unwrap(), 0);
4977    }
4978}
4979
4980/// Blue/green swap slots: a candidate registered beside the active endpoint,
4981/// promoted by `cutover_candidate`, with the old incumbent drained by endpoint.
4982#[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        // Every by-id consumer still resolves the incumbent.
5062        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        // Connection-keyed lookups see the candidate, so its own frames resolve.
5081        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        // One candidate per id.
5096        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    /// The cutover linearization point. A relay reserved on the incumbent before
5112    /// cutover must never become a route on the incumbent, and every relay
5113    /// reserved after it must land on the promoted candidate.
5114    #[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        // A route.open reserved after cutover goes to the promoted candidate.
5137        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        // The incumbent acks the early relay after cutover.
5147        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        // The late relay commits on the candidate; only its pair stays reserved
5178        // until then.
5179        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        // One bound route and one in-flight relay on the incumbent.
5196        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    /// An endpoint replaced WITHOUT a promotion (a successor registered over it
5253    /// as an ordinary active HELLO) is stale, not superseded, and its ack still
5254    /// fails the acking connection exactly as it did before swap slots existed:
5255    /// `StaleModuleEndpoint`, the reservation indexes already stripped, and the
5256    /// waiting client's sender dropped unanswered.
5257    #[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        // A candidate whose connection drops leaves the incumbent routable.
5307        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        // After a cutover, the incumbent's teardown releases its own routes and
5324        // leaves the promoted candidate in place.
5325        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}