Skip to main content

fips_core/node/handlers/
rx_loop.rs

1//! RX event loop and dataplane dispatch.
2
3use crate::control::queries;
4use crate::control::{ControlMessage, ControlSenders, ControlSocket, commands};
5use crate::dataplane::DataplaneFastIngressRx;
6use crate::node::{
7    EndpointDataBatchRx, EndpointEventSender, Node, NodeError, endpoint_data_batch_channel,
8    lifecycle::NetworkRebindCompletion,
9};
10use crate::transport::PacketRx;
11use crate::upper::tun::TunOutboundRx;
12use std::sync::Arc;
13use std::time::{Duration, Instant};
14use tokio::sync::Notify;
15use tokio::sync::mpsc::Receiver;
16use tracing::{debug, info, warn};
17
18mod budget;
19mod dataplane;
20mod drain;
21
22#[cfg(test)]
23mod tests;
24
25use budget::*;
26use drain::*;
27
28pub(in crate::node) struct RxLoopDataplaneIo<'a> {
29    packet_rx: &'a mut PacketRx,
30    dataplane_fast_ingress_rx: &'a mut DataplaneFastIngressRx,
31    endpoint_data_rx: &'a mut EndpointDataBatchRx,
32    tun_outbound_rx: &'a mut TunOutboundRx,
33    endpoint_tx: &'a EndpointEventSender,
34}
35
36struct RxLoopDataplaneRuntime {
37    packet_rx: PacketRx,
38    dataplane_fast_ingress_rx: DataplaneFastIngressRx,
39    endpoint_data_rx: EndpointDataBatchRx,
40    tun_outbound_rx: TunOutboundRx,
41    endpoint_tx: EndpointEventSender,
42}
43
44impl RxLoopDataplaneRuntime {
45    fn io(&mut self) -> RxLoopDataplaneIo<'_> {
46        RxLoopDataplaneIo {
47            packet_rx: &mut self.packet_rx,
48            dataplane_fast_ingress_rx: &mut self.dataplane_fast_ingress_rx,
49            endpoint_data_rx: &mut self.endpoint_data_rx,
50            tun_outbound_rx: &mut self.tun_outbound_rx,
51            endpoint_tx: &self.endpoint_tx,
52        }
53    }
54}
55
56#[derive(Clone, Copy, Debug, Eq, PartialEq)]
57pub(in crate::node) struct RxLoopDataplaneTurnLimits {
58    packet: usize,
59    endpoint: usize,
60    tun: usize,
61    crypto: usize,
62}
63
64impl RxLoopDataplaneTurnLimits {
65    pub(in crate::node) fn new(packet: usize, endpoint: usize, tun: usize, crypto: usize) -> Self {
66        Self {
67            packet,
68            endpoint,
69            tun,
70            crypto,
71        }
72    }
73}
74
75#[cfg(test)]
76pub(in crate::node) fn rx_loop_dataplane_io<'a>(
77    packet_rx: &'a mut PacketRx,
78    dataplane_fast_ingress_rx: &'a mut DataplaneFastIngressRx,
79    endpoint_data_rx: &'a mut EndpointDataBatchRx,
80    tun_outbound_rx: &'a mut TunOutboundRx,
81    endpoint_tx: &'a EndpointEventSender,
82) -> RxLoopDataplaneIo<'a> {
83    RxLoopDataplaneIo {
84        packet_rx,
85        dataplane_fast_ingress_rx,
86        endpoint_data_rx,
87        tun_outbound_rx,
88        endpoint_tx,
89    }
90}
91
92impl Node {
93    /// Run the receive event loop.
94    ///
95    /// Processes packets from all transports, dispatching based on
96    /// the phase field in the 4-byte common prefix:
97    /// - Phase 0x0: Encrypted frame (session data)
98    /// - Phase 0x1: Handshake message 1 (initiator -> responder)
99    /// - Phase 0x2: Handshake message 2 (responder -> initiator)
100    ///
101    /// Also processes outbound IPv6 packets from the TUN reader for session
102    /// encapsulation and routing through the mesh.
103    ///
104    /// Also processes DNS-resolved identities for identity cache population.
105    ///
106    /// Also runs a periodic tick (1s) to clean up stale handshake connections
107    /// that never received a response. This prevents resource leaks when peers
108    /// are unreachable.
109    ///
110    /// This method takes ownership of the packet_rx channel and runs
111    /// until the channel is closed (typically when stop() is called).
112    pub async fn run_rx_loop(&mut self) -> Result<(), NodeError> {
113        let packet_rx = self.packet_rx.take().ok_or(NodeError::NotStarted)?;
114
115        // Take the TUN outbound receiver, or create a dummy channel that never
116        // produces messages (when TUN is disabled). Holding the sender prevents
117        // the channel from closing.
118        let (tun_outbound_rx, _tun_guard) = match self.tun_outbound_rx.take() {
119            Some(rx) => (rx, None),
120            None => {
121                let (tx, rx) = crate::upper::tun::tun_outbound_channel(1);
122                (rx, Some(tx))
123            }
124        };
125
126        // Take the DNS identity receiver, or create a dummy channel (when DNS
127        // is disabled). Same pattern as TUN outbound.
128        let (mut dns_identity_rx, _dns_guard) = match self.dns_identity_rx.take() {
129            Some(rx) => (rx, None),
130            None => {
131                let (tx, rx) = tokio::sync::mpsc::channel(1);
132                (rx, Some(tx))
133            }
134        };
135
136        // Take the endpoint control receiver, or create a dummy channel
137        // when the embedded endpoint API is not in use.
138        let (mut endpoint_control_rx, _endpoint_control_guard) =
139            match self.endpoint_control_rx.take() {
140                Some(rx) => (rx, None),
141                None => {
142                    let (tx, rx) = tokio::sync::mpsc::channel(1);
143                    (rx, Some(tx))
144                }
145            };
146        let (endpoint_data_rx, _endpoint_data_guard) = match self.endpoint_data_rx.take() {
147            Some(rx) => (rx, None),
148            None => {
149                let (tx, rx) = endpoint_data_batch_channel(1);
150                (rx, Some(tx))
151            }
152        };
153        let (dataplane_fast_ingress_rx, _dataplane_fast_ingress_guard) =
154            match self.dataplane_fast_ingress_rx.take() {
155                Some(rx) => (rx, None),
156                None => {
157                    let (tx, rx) = tokio::sync::mpsc::channel(1);
158                    (rx, Some(tx))
159                }
160            };
161
162        let mut tick =
163            tokio::time::interval(Duration::from_secs(self.config.node.tick_interval_secs));
164        tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
165        let mut maintenance_state = RxLoopMaintenanceState::default();
166        let (network_rebind_completion_tx, mut network_rebind_completion_rx) =
167            tokio::sync::mpsc::channel::<NetworkRebindCompletion>(1);
168        let mut network_rebind_in_progress = false;
169
170        // Set up control socket channels. Read-only queries are separated
171        // from mutating commands so operator status reads can get reserved
172        // progress while a command awaits slower discovery/transport work.
173        let (control_query_tx, mut control_query_rx) =
174            tokio::sync::mpsc::channel::<ControlMessage>(32);
175        let (control_command_tx, mut control_command_rx) =
176            tokio::sync::mpsc::channel::<ControlMessage>(32);
177
178        if self.config.node.control.enabled {
179            let config = self.config.node.control.clone();
180            let senders = ControlSenders::new(control_query_tx.clone(), control_command_tx.clone());
181            tokio::spawn(async move {
182                match ControlSocket::bind(&config) {
183                    Ok(socket) => {
184                        socket.accept_loop_split(senders).await;
185                    }
186                    Err(e) => {
187                        warn!(error = %e, "Failed to bind control socket");
188                    }
189                }
190            });
191        }
192        // Drop unused sender to avoid keeping channel open if control is disabled
193        drop(control_query_tx);
194        drop(control_command_tx);
195
196        let dataplane_endpoint_tx = self.endpoint_events.sender().unwrap_or_else(|| {
197            let (tx, rx) = EndpointEventSender::channel(1);
198            drop(rx);
199            tx
200        });
201        let dataplane_readiness_notify = self.dataplane.readiness_notify();
202        let nostr_node_event_notify = self
203            .nostr_discovery
204            .as_ref()
205            .map(|discovery| discovery.node_event_notify());
206        let mut dataplane_runtime = RxLoopDataplaneRuntime {
207            packet_rx,
208            dataplane_fast_ingress_rx,
209            endpoint_data_rx,
210            tun_outbound_rx,
211            endpoint_tx: dataplane_endpoint_tx,
212        };
213
214        info!("RX event loop started");
215        // Optional perf profiler (FIPS_PERF=1). No-op otherwise.
216        crate::perf_profile::maybe_spawn_reporter();
217        // Tokio intervals tick immediately on first poll. Consume that startup
218        // tick so the reserved-progress branch below represents a due periodic
219        // maintenance turn, not an eager pre-data maintenance pass.
220        tick.tick().await;
221        let mut nostr_event_turn_not_before = Instant::now();
222
223        loop {
224            tokio::select! {
225                biased;
226                // Timer-driven liveness is a reserved-progress branch. It
227                // performs bounded pre/post data drains and timeboxes slow
228                // discovery/status work, so hot packet or endpoint/TUN queues
229                // cannot indefinitely postpone heartbeat, rekey, MMP, route
230                // aging, or path maintenance.
231                _ = tick.tick() => {
232                    let drained = {
233                        let mut dataplane_io = dataplane_runtime.io();
234                        self.drain_rx_loop_data_queues(
235                            &mut dataplane_io,
236                            ENDPOINT_DRAIN_BUDGET,
237                        ).await
238                    };
239                    if drained.has_drained() {
240                        maintenance_state.record_data_activity(Instant::now());
241                        debug!(
242                            drained = drained.total(),
243                            drained_packets = drained.packets,
244                            drained_tun = drained.tun,
245                            drained_endpoint = drained.endpoint,
246                            "Drained queued packets before rx-loop maintenance"
247                        );
248                    }
249                    let maintenance_plan = maintenance_state.plan_maintenance(
250                        drained,
251                        Instant::now(),
252                        RX_LOOP_RECENT_DATA_ACTIVITY_WINDOW,
253                        RX_LOOP_SLOW_MAINTENANCE_IDLE_TIMEOUT,
254                        RX_LOOP_SLOW_MAINTENANCE_BUSY_TIMEOUT,
255                    );
256
257                    let slow_timed_out = self.run_rx_loop_maintenance_tick(
258                        maintenance_plan,
259                    ).await;
260                    maintenance_state.record_maintenance_result(
261                        maintenance_plan.data_pressure(),
262                        slow_timed_out,
263                    );
264
265                    let post_drained = {
266                        let mut dataplane_io = dataplane_runtime.io();
267                        self.drain_rx_loop_data_queues(
268                            &mut dataplane_io,
269                            PACKET_DRAIN_BUDGET,
270                        ).await
271                    };
272                    if post_drained.has_drained() {
273                        maintenance_state.record_data_activity(Instant::now());
274                        debug!(
275                            drained = post_drained.total(),
276                            drained_packets = post_drained.packets,
277                            drained_tun = post_drained.tun,
278                            drained_endpoint = post_drained.endpoint,
279                            "Drained queued packets after rx-loop maintenance"
280                        );
281                    }
282                }
283                Some(message) = control_query_rx.recv() => {
284                    self.drain_control_queries(
285                        &mut control_query_rx,
286                        Some(message),
287                        ENDPOINT_DRAIN_BUDGET,
288                    ).await;
289                }
290                Some(completion) = network_rebind_completion_rx.recv() => {
291                    network_rebind_in_progress = false;
292                    self.complete_network_rebind(completion).await;
293                }
294                // Endpoint control carries management/lifecycle commands.
295                // Endpoint payload batches stay on the data lane; this branch
296                // keeps control work from waiting behind hot raw receive.
297                // Endpoint data batches intentionally remain below packet_rx.
298                Some(command) = endpoint_control_rx.recv() => {
299                    if let Some(request) = self.handle_endpoint_control(command).await {
300                        if network_rebind_in_progress {
301                            request.reject(NodeError::TransportError(
302                                "network transport rebind already in progress".to_string(),
303                            ));
304                        } else {
305                            network_rebind_in_progress = true;
306                            self.spawn_network_rebind_preparation(
307                                request,
308                                network_rebind_completion_tx.clone(),
309                            );
310                        }
311                    }
312                }
313                // Discovery receives an explicitly rate-limited fair turn:
314                // lifecycle control stays reserved, while hot discovery can
315                // delay dataplane work by at most one timeboxed turn per gap.
316                _ = wait_for_optional_notify_after(
317                    nostr_node_event_notify.as_ref(),
318                    nostr_event_turn_not_before,
319                ) => {
320                    self.poll_nostr_discovery_event_turn(
321                        RX_LOOP_NOSTR_EVENT_TURN_BUDGET,
322                    ).await;
323                    nostr_event_turn_not_before =
324                        Instant::now() + RX_LOOP_NOSTR_EVENT_TURN_INTERVAL;
325                }
326                packet = dataplane_runtime.packet_rx.recv() => {
327                    match packet {
328                        Some(p) => {
329                            let latency_packet = p.is_transport_priority();
330                            let mut firsts = crate::dataplane::DataplaneLiveTurnFirsts {
331                                raw_packet: Some(p),
332                                ..Default::default()
333                            };
334                            if let Ok(packet) = dataplane_runtime.tun_outbound_rx.try_recv() {
335                                firsts.tun_packet = Some(packet);
336                            }
337                            let latency_work_ready = latency_packet
338                                || dataplane_runtime.packet_rx.priority_ready_packets() > 0;
339                            if latency_work_ready {
340                                let packet_budget = packet_drain_budget(true);
341                                let endpoint_budget = endpoint_drain_budget(packet_budget);
342                                let tun_budget = tun_drain_budget(packet_budget);
343                                let crypto_budget = mixed_dataplane_crypto_budget(
344                                    packet_budget,
345                                    endpoint_budget,
346                                    tun_budget,
347                                );
348                                let mut turn = {
349                                    let mut dataplane_io = dataplane_runtime.io();
350                                    self.drain_dataplane_turn_with_firsts(
351                                        &mut dataplane_io,
352                                        firsts,
353                                        RxLoopDataplaneTurnLimits::new(
354                                            packet_budget,
355                                            endpoint_budget,
356                                            tun_budget,
357                                            crypto_budget,
358                                        ),
359                                    ).await
360                                };
361                                self.finish_dataplane_turn(
362                                    &mut turn,
363                                    &mut maintenance_state,
364                                    &mut control_query_rx,
365                                    CONTROL_QUERY_INTERLEAVE_BUDGET,
366                                ).await;
367                            } else {
368                                firsts.raw_ingress_prefetch = true;
369                                let mut dataplane_io = dataplane_runtime.io();
370                                self.service_dataplane_bulk_turns(
371                                    &mut dataplane_io,
372                                    firsts,
373                                    &mut maintenance_state,
374                                    &mut control_query_rx,
375                                ).await;
376                            }
377                        }
378                        None => break, // channel closed
379                    }
380                }
381                Some(fast_ingress) = dataplane_runtime.dataplane_fast_ingress_rx.recv() => {
382                    let mut dataplane_io = dataplane_runtime.io();
383                    self.service_dataplane_bulk_turns(
384                        &mut dataplane_io,
385                        crate::dataplane::DataplaneLiveTurnFirsts {
386                            fast_ingress: Some(fast_ingress),
387                            ..Default::default()
388                        },
389                        &mut maintenance_state,
390                        &mut control_query_rx,
391                    ).await;
392                }
393                _ = dataplane_readiness_notify.notified() => {
394                    Box::pin(self.flush_pending_local_rendezvous_sync()).await;
395                    let mut dataplane_io = dataplane_runtime.io();
396                    self.service_dataplane_completion_turns(
397                        &mut dataplane_io,
398                        &mut maintenance_state,
399                        &mut control_query_rx,
400                    ).await;
401                }
402                Some(ipv6_packet) = dataplane_runtime.tun_outbound_rx.recv() => {
403                    let tun_budget = tun_drain_budget(LATENCY_PACKET_DRAIN_BUDGET);
404                    let mut turn = {
405                        let mut dataplane_io = dataplane_runtime.io();
406                        self.drain_dataplane_turn_with_firsts(
407                            &mut dataplane_io,
408                            crate::dataplane::DataplaneLiveTurnFirsts {
409                                tun_packet: Some(ipv6_packet),
410                                ..Default::default()
411                            },
412                            RxLoopDataplaneTurnLimits::new(0, 0, tun_budget, tun_budget),
413                        ).await
414                    };
415                    self.finish_dataplane_turn(
416                        &mut turn,
417                        &mut maintenance_state,
418                        &mut control_query_rx,
419                        0,
420                    ).await;
421                }
422                Some(identity) = dns_identity_rx.recv() => {
423                    debug!(
424                        node_addr = %identity.node_addr,
425                        "Registering identity from DNS resolution"
426                    );
427                    self.register_dns_identity(identity.node_addr, identity.pubkey);
428                }
429                Some(batch) = dataplane_runtime.endpoint_data_rx.recv() => {
430                    let mut turn = {
431                        let mut dataplane_io = dataplane_runtime.io();
432                        self.drain_dataplane_turn_with_firsts(
433                            &mut dataplane_io,
434                            crate::dataplane::DataplaneLiveTurnFirsts {
435                                endpoint_data_batch: Some(batch),
436                                ..Default::default()
437                            },
438                            RxLoopDataplaneTurnLimits::new(
439                                0,
440                                ENDPOINT_DRAIN_BUDGET,
441                                0,
442                                PACKET_DRAIN_BUDGET,
443                            ),
444                        ).await
445                    };
446                    self.finish_dataplane_turn(
447                        &mut turn,
448                        &mut maintenance_state,
449                        &mut control_query_rx,
450                        0,
451                    ).await;
452                }
453                Some((request, response_tx)) = control_command_rx.recv() => {
454                    let response = commands::dispatch(
455                        self,
456                        &request.command,
457                        request.params.as_ref(),
458                    ).await;
459                    let _ = response_tx.send(response);
460                }
461            }
462        }
463
464        info!("RX event loop stopped (channel closed)");
465        Ok(())
466    }
467
468    async fn drain_rx_loop_data_queues(
469        &mut self,
470        io: &mut RxLoopDataplaneIo<'_>,
471        budget: usize,
472    ) -> RxLoopDataDrainStats {
473        let fast_ingress =
474            Self::take_dataplane_fast_ingress_batch(io.dataplane_fast_ingress_rx, budget);
475        let packet_budget = budget.max(
476            fast_ingress
477                .as_ref()
478                .map_or(0, |fast_ingress| fast_ingress.len()),
479        );
480        let endpoint_budget = endpoint_drain_budget(packet_budget);
481        let tun_budget = tun_drain_budget(packet_budget);
482        let crypto_budget =
483            mixed_dataplane_crypto_budget(packet_budget, endpoint_budget, tun_budget);
484        let mut turn = self
485            .drain_dataplane_turn_with_firsts(
486                io,
487                crate::dataplane::DataplaneLiveTurnFirsts {
488                    fast_ingress,
489                    ..Default::default()
490                },
491                RxLoopDataplaneTurnLimits::new(
492                    packet_budget,
493                    endpoint_budget,
494                    tun_budget,
495                    crypto_budget,
496                ),
497            )
498            .await;
499        let drained_packets = Self::dataplane_packet_activity(&turn);
500        let control_drained = Box::pin(self.process_dataplane_control_ingress(&mut turn)).await;
501        RxLoopDataDrainStats::new(
502            drained_packets,
503            turn.tun_source_drained(),
504            turn.endpoint_source_drained(),
505            control_drained,
506        )
507    }
508
509    fn take_dataplane_fast_ingress_batch(
510        dataplane_fast_ingress_rx: &mut crate::dataplane::DataplaneFastIngressRx,
511        limit: usize,
512    ) -> Option<crate::dataplane::DataplaneFastIngressBatch> {
513        let fast_ingress = dataplane_fast_ingress_rx.try_recv().ok()?;
514        Some(Self::coalesce_dataplane_fast_ingress(
515            fast_ingress,
516            dataplane_fast_ingress_rx,
517            limit,
518        ))
519    }
520
521    fn coalesce_dataplane_fast_ingress(
522        mut fast_ingress: crate::dataplane::DataplaneFastIngressBatch,
523        dataplane_fast_ingress_rx: &mut crate::dataplane::DataplaneFastIngressRx,
524        limit: usize,
525    ) -> crate::dataplane::DataplaneFastIngressBatch {
526        while fast_ingress.len() < limit {
527            let Ok(next) = dataplane_fast_ingress_rx.try_recv() else {
528                break;
529            };
530            fast_ingress.absorb(next);
531        }
532        fast_ingress
533    }
534
535    async fn service_dataplane_bulk_turns(
536        &mut self,
537        io: &mut RxLoopDataplaneIo<'_>,
538        firsts: crate::dataplane::DataplaneLiveTurnFirsts,
539        maintenance_state: &mut RxLoopMaintenanceState,
540        control_query_rx: &mut Receiver<ControlMessage>,
541    ) {
542        let started = Instant::now();
543        let mut firsts = Some(firsts);
544        let mut turns = 0usize;
545
546        loop {
547            if turns > 0
548                && (turns >= RX_LOOP_BULK_SERVICE_MAX_TURNS
549                    || started.elapsed() >= RX_LOOP_BULK_SERVICE_MAX_ELAPSED
550                    || io.packet_rx.priority_ready_packets() > 0)
551            {
552                break;
553            }
554
555            let packet_budget = PACKET_DRAIN_BUDGET;
556            let mut turn_firsts = firsts.take().unwrap_or_default();
557            turn_firsts.raw_ingress_prefetch = true;
558            turn_firsts.fast_ingress = match turn_firsts.fast_ingress.take() {
559                Some(fast_ingress) => Some(Self::coalesce_dataplane_fast_ingress(
560                    fast_ingress,
561                    io.dataplane_fast_ingress_rx,
562                    packet_budget,
563                )),
564                None => Self::take_dataplane_fast_ingress_batch(
565                    io.dataplane_fast_ingress_rx,
566                    packet_budget,
567                ),
568            };
569            let packet_budget = packet_budget.max(
570                turn_firsts
571                    .fast_ingress
572                    .as_ref()
573                    .map_or(0, |fast_ingress| fast_ingress.len()),
574            );
575            let endpoint_budget = endpoint_drain_budget(packet_budget);
576            let tun_budget = tun_drain_budget(packet_budget);
577            let crypto_budget =
578                mixed_dataplane_crypto_budget(packet_budget, endpoint_budget, tun_budget);
579
580            let mut turn = self
581                .drain_dataplane_turn_with_firsts(
582                    io,
583                    turn_firsts,
584                    RxLoopDataplaneTurnLimits::new(
585                        packet_budget,
586                        endpoint_budget,
587                        tun_budget,
588                        crypto_budget,
589                    ),
590                )
591                .await;
592            let raw_drained = Self::dataplane_raw_ingress_activity(&turn);
593            let control_activity = Self::dataplane_control_activity(&turn);
594            let completions_drained = turn.summary().completions();
595            let admission_dropped =
596                turn.summary().inbound_dropped() > 0 || turn.summary().outbound_dropped() > 0;
597            let keep_servicing = !admission_dropped
598                && (raw_drained >= packet_budget
599                    || completions_drained >= crypto_budget
600                    || turn.tun_source_drained() >= tun_budget
601                    || turn.endpoint_source_drained() >= endpoint_budget);
602            let control_drained = self
603                .finish_dataplane_turn(
604                    &mut turn,
605                    maintenance_state,
606                    control_query_rx,
607                    CONTROL_QUERY_INTERLEAVE_BUDGET,
608                )
609                .await;
610            turns += 1;
611            let mut runnable_work = self.dataplane.has_runnable_work();
612
613            if control_drained == 0
614                && bulk_admission_pressure_relief_due(
615                    admission_dropped,
616                    runnable_work,
617                    turns,
618                    started.elapsed(),
619                    io.packet_rx.priority_ready_packets(),
620                )
621            {
622                let mut relief_turn = self
623                    .drain_dataplane_completion_turn(io, LATENCY_PACKET_DRAIN_BUDGET)
624                    .await;
625                let relief_control_drained = self
626                    .finish_dataplane_turn(&mut relief_turn, maintenance_state, control_query_rx, 0)
627                    .await;
628                turns += 1;
629                runnable_work = self.dataplane.has_runnable_work();
630                if relief_control_drained > 0 || !runnable_work {
631                    break;
632                }
633            }
634
635            if control_drained > 0
636                || admission_dropped
637                || (!keep_servicing && !runnable_work)
638                || (control_activity > 0 && !runnable_work)
639            {
640                break;
641            }
642        }
643    }
644
645    async fn service_dataplane_completion_turns(
646        &mut self,
647        io: &mut RxLoopDataplaneIo<'_>,
648        maintenance_state: &mut RxLoopMaintenanceState,
649        control_query_rx: &mut Receiver<ControlMessage>,
650    ) {
651        let started = Instant::now();
652        let mut turns = 0usize;
653
654        loop {
655            if turns > 0
656                && (turns >= RX_LOOP_BULK_SERVICE_MAX_TURNS
657                    || started.elapsed() >= RX_LOOP_BULK_SERVICE_MAX_ELAPSED
658                    || io.packet_rx.priority_ready_packets() > 0)
659            {
660                break;
661            }
662
663            let mut turn = self
664                .drain_dataplane_completion_turn(io, LATENCY_PACKET_DRAIN_BUDGET)
665                .await;
666            let control_drained = self
667                .finish_dataplane_turn(&mut turn, maintenance_state, control_query_rx, 0)
668                .await;
669            turns += 1;
670
671            let runnable_work = self.dataplane.has_runnable_work();
672            if control_drained > 0 || !runnable_work {
673                break;
674            }
675        }
676    }
677
678    async fn finish_dataplane_turn(
679        &mut self,
680        turn: &mut crate::dataplane::DataplaneLiveNodeTurn,
681        maintenance_state: &mut RxLoopMaintenanceState,
682        control_query_rx: &mut Receiver<ControlMessage>,
683        control_query_budget: usize,
684    ) -> usize {
685        let had_activity = turn.has_activity();
686        let control_drained = Box::pin(self.process_dataplane_control_ingress(turn))
687            .await
688            .saturating_add(Box::pin(self.drain_deferred_dataplane_control_turns()).await);
689        if control_drained > 0 && self.dataplane.has_deferred_raw_ingress() {
690            self.dataplane.readiness_notify().notify_one();
691        }
692        let query_drained = if control_query_budget > 0 {
693            self.drain_control_queries(control_query_rx, None, control_query_budget)
694                .await
695        } else {
696            0
697        };
698        if had_activity || control_drained > 0 {
699            maintenance_state.record_data_activity(Instant::now());
700        }
701        control_drained.saturating_add(query_drained)
702    }
703
704    async fn drain_control_queries(
705        &mut self,
706        control_query_rx: &mut Receiver<ControlMessage>,
707        first_message: Option<ControlMessage>,
708        budget: usize,
709    ) -> usize {
710        let mut drain = SingleLaneDrainCursor::new(first_message, budget);
711        while let Some((request, response_tx)) = drain.next(control_query_rx) {
712            let response = queries::dispatch(self, &request.command, request.params.as_ref());
713            let _ = response_tx.send(response);
714        }
715
716        drain.drained()
717    }
718
719    async fn run_rx_loop_maintenance_tick(&mut self, plan: RxLoopMaintenancePlan) -> bool {
720        if !rx_loop_fast_maintenance_within_budget(self.run_rx_loop_fast_maintenance_tick()).await {
721            crate::perf_profile::record_event(
722                crate::perf_profile::Event::RxLoopSlowMaintenanceTimeout,
723            );
724            self.mark_rx_loop_maintenance_timeout();
725            warn!(
726                timeout_ms = RX_LOOP_FAST_MAINTENANCE_TIMEOUT.as_millis() as u64,
727                data_pressure = plan.data_pressure(),
728                "RX loop liveness maintenance timed out; continuing packet processing"
729            );
730            return true;
731        }
732
733        let Some(slow_timeout) = plan.slow_timeout() else {
734            crate::perf_profile::record_event(
735                crate::perf_profile::Event::RxLoopSlowMaintenanceSkipped,
736            );
737            return false;
738        };
739
740        if tokio::time::timeout(slow_timeout, self.run_rx_loop_slow_maintenance_tick())
741            .await
742            .is_err()
743        {
744            crate::perf_profile::record_event(
745                crate::perf_profile::Event::RxLoopSlowMaintenanceTimeout,
746            );
747            self.mark_rx_loop_maintenance_timeout();
748            warn!(
749                timeout_ms = slow_timeout.as_millis() as u64,
750                data_pressure = plan.data_pressure(),
751                "RX loop slow maintenance timed out; continuing packet processing"
752            );
753            return true;
754        }
755        false
756    }
757
758    async fn run_rx_loop_fast_maintenance_tick(&mut self) {
759        self.check_timeouts();
760        let now_ms = Self::now_ms();
761        // Link/session liveness must run before slower retry/discovery work:
762        // under bulk send pressure a late heartbeat or MMP report is
763        // indistinguishable from a dead direct path on the remote peer.
764        self.check_link_heartbeats().await;
765        self.reload_peer_acl();
766        self.resend_pending_handshakes(now_ms).await;
767        self.resend_pending_rekeys(now_ms).await;
768        self.resend_pending_session_handshakes(now_ms).await;
769        self.resend_pending_session_msg3(now_ms).await;
770        self.retry_pending_session_traffic().await;
771        self.purge_idle_sessions(now_ms);
772        self.purge_learned_routes(now_ms);
773        self.check_mmp_reports().await;
774        self.check_session_mmp_reports().await;
775        self.check_rekey().await;
776        self.check_session_rekey().await;
777        self.check_pending_lookups(now_ms).await;
778        self.poll_pending_connects().await;
779        self.process_pending_retries(now_ms).await;
780        self.poll_transport_discovery().await;
781        self.sample_transport_congestion();
782    }
783
784    async fn run_rx_loop_slow_maintenance_tick(&mut self) {
785        if let Some(delay) = rx_loop_slow_maintenance_fault_delay() {
786            tokio::time::sleep(delay).await;
787        }
788
789        // Discovery and graph/stat maintenance can involve relay work or
790        // larger scans. Keep it bounded after direct-path liveness and session
791        // upkeep so a slow Nostr/LAN tick degrades discovery freshness, not
792        // packet flow.
793        self.poll_nostr_discovery().await;
794        self.poll_lan_discovery().await;
795        self.poll_local_rendezvous().await;
796        self.check_tree_state().await;
797        self.check_bloom_state().await;
798        self.compute_mesh_size();
799        self.record_stats_history();
800    }
801}
802
803async fn wait_for_optional_notify(notify: Option<&Arc<Notify>>) {
804    match notify {
805        Some(notify) => notify.notified().await,
806        None => std::future::pending().await,
807    }
808}
809
810async fn wait_for_optional_notify_after(notify: Option<&Arc<Notify>>, not_before: Instant) {
811    tokio::time::sleep(not_before.saturating_duration_since(Instant::now())).await;
812    wait_for_optional_notify(notify).await;
813}
814
815async fn rx_loop_fast_maintenance_within_budget<F>(maintenance: F) -> bool
816where
817    F: std::future::Future<Output = ()>,
818{
819    tokio::time::timeout(RX_LOOP_FAST_MAINTENANCE_TIMEOUT, maintenance)
820        .await
821        .is_ok()
822}
823
824fn bulk_admission_pressure_relief_due(
825    admission_dropped: bool,
826    runnable_work: bool,
827    turns: usize,
828    elapsed: Duration,
829    priority_ready_packets: usize,
830) -> bool {
831    admission_dropped
832        && runnable_work
833        && turns < RX_LOOP_BULK_SERVICE_MAX_TURNS
834        && elapsed < RX_LOOP_BULK_SERVICE_MAX_ELAPSED
835        && priority_ready_packets == 0
836}