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
222        loop {
223            tokio::select! {
224                biased;
225                // Timer-driven liveness is a reserved-progress branch. It
226                // performs bounded pre/post data drains and timeboxes slow
227                // discovery/status work, so hot packet or endpoint/TUN queues
228                // cannot indefinitely postpone heartbeat, rekey, MMP, route
229                // aging, or path maintenance.
230                _ = tick.tick() => {
231                    let drained = {
232                        let mut dataplane_io = dataplane_runtime.io();
233                        self.drain_rx_loop_data_queues(
234                            &mut dataplane_io,
235                            ENDPOINT_DRAIN_BUDGET,
236                        ).await
237                    };
238                    if drained.has_drained() {
239                        maintenance_state.record_data_activity(Instant::now());
240                        debug!(
241                            drained = drained.total(),
242                            drained_packets = drained.packets,
243                            drained_tun = drained.tun,
244                            drained_endpoint = drained.endpoint,
245                            "Drained queued packets before rx-loop maintenance"
246                        );
247                    }
248                    let maintenance_plan = maintenance_state.plan_maintenance(
249                        drained,
250                        Instant::now(),
251                        RX_LOOP_RECENT_DATA_ACTIVITY_WINDOW,
252                        RX_LOOP_SLOW_MAINTENANCE_IDLE_TIMEOUT,
253                        RX_LOOP_SLOW_MAINTENANCE_BUSY_TIMEOUT,
254                    );
255
256                    let slow_timed_out = self.run_rx_loop_maintenance_tick(
257                        maintenance_plan,
258                    ).await;
259                    maintenance_state.record_maintenance_result(
260                        maintenance_plan.data_pressure(),
261                        slow_timed_out,
262                    );
263
264                    let post_drained = {
265                        let mut dataplane_io = dataplane_runtime.io();
266                        self.drain_rx_loop_data_queues(
267                            &mut dataplane_io,
268                            PACKET_DRAIN_BUDGET,
269                        ).await
270                    };
271                    if post_drained.has_drained() {
272                        maintenance_state.record_data_activity(Instant::now());
273                        debug!(
274                            drained = post_drained.total(),
275                            drained_packets = post_drained.packets,
276                            drained_tun = post_drained.tun,
277                            drained_endpoint = post_drained.endpoint,
278                            "Drained queued packets after rx-loop maintenance"
279                        );
280                    }
281                }
282                _ = wait_for_optional_notify(nostr_node_event_notify.as_ref()) => {
283                    self.poll_nostr_discovery().await;
284                }
285                Some(message) = control_query_rx.recv() => {
286                    self.drain_control_queries(
287                        &mut control_query_rx,
288                        Some(message),
289                        ENDPOINT_DRAIN_BUDGET,
290                    ).await;
291                }
292                Some(completion) = network_rebind_completion_rx.recv() => {
293                    network_rebind_in_progress = false;
294                    self.complete_network_rebind(completion).await;
295                }
296                // Endpoint control carries management/lifecycle commands.
297                // Endpoint payload batches stay on the data lane; this branch
298                // keeps control work from waiting behind hot raw receive.
299                // Endpoint data batches intentionally remain below packet_rx.
300                Some(command) = endpoint_control_rx.recv() => {
301                    if let Some(request) = self.handle_endpoint_control(command).await {
302                        if network_rebind_in_progress {
303                            request.reject(NodeError::TransportError(
304                                "network transport rebind already in progress".to_string(),
305                            ));
306                        } else {
307                            network_rebind_in_progress = true;
308                            self.spawn_network_rebind_preparation(
309                                request,
310                                network_rebind_completion_tx.clone(),
311                            );
312                        }
313                    }
314                }
315                packet = dataplane_runtime.packet_rx.recv() => {
316                    match packet {
317                        Some(p) => {
318                            let latency_packet = p.is_transport_priority();
319                            let mut firsts = crate::dataplane::DataplaneLiveTurnFirsts {
320                                raw_packet: Some(p),
321                                ..Default::default()
322                            };
323                            if let Ok(packet) = dataplane_runtime.tun_outbound_rx.try_recv() {
324                                firsts.tun_packet = Some(packet);
325                            }
326                            let latency_work_ready = latency_packet
327                                || dataplane_runtime.packet_rx.priority_ready_packets() > 0;
328                            if latency_work_ready {
329                                let packet_budget = packet_drain_budget(true);
330                                let endpoint_budget = endpoint_drain_budget(packet_budget);
331                                let tun_budget = tun_drain_budget(packet_budget);
332                                let crypto_budget = mixed_dataplane_crypto_budget(
333                                    packet_budget,
334                                    endpoint_budget,
335                                    tun_budget,
336                                );
337                                let mut turn = {
338                                    let mut dataplane_io = dataplane_runtime.io();
339                                    self.drain_dataplane_turn_with_firsts(
340                                        &mut dataplane_io,
341                                        firsts,
342                                        RxLoopDataplaneTurnLimits::new(
343                                            packet_budget,
344                                            endpoint_budget,
345                                            tun_budget,
346                                            crypto_budget,
347                                        ),
348                                    ).await
349                                };
350                                self.finish_dataplane_turn(
351                                    &mut turn,
352                                    &mut maintenance_state,
353                                    &mut control_query_rx,
354                                    CONTROL_QUERY_INTERLEAVE_BUDGET,
355                                ).await;
356                            } else {
357                                firsts.raw_ingress_prefetch = true;
358                                let mut dataplane_io = dataplane_runtime.io();
359                                self.service_dataplane_bulk_turns(
360                                    &mut dataplane_io,
361                                    firsts,
362                                    &mut maintenance_state,
363                                    &mut control_query_rx,
364                                ).await;
365                            }
366                        }
367                        None => break, // channel closed
368                    }
369                }
370                Some(fast_ingress) = dataplane_runtime.dataplane_fast_ingress_rx.recv() => {
371                    let mut dataplane_io = dataplane_runtime.io();
372                    self.service_dataplane_bulk_turns(
373                        &mut dataplane_io,
374                        crate::dataplane::DataplaneLiveTurnFirsts {
375                            fast_ingress: Some(fast_ingress),
376                            ..Default::default()
377                        },
378                        &mut maintenance_state,
379                        &mut control_query_rx,
380                    ).await;
381                }
382                _ = dataplane_readiness_notify.notified() => {
383                    Box::pin(self.flush_pending_local_rendezvous_sync()).await;
384                    let mut dataplane_io = dataplane_runtime.io();
385                    self.service_dataplane_completion_turns(
386                        &mut dataplane_io,
387                        &mut maintenance_state,
388                        &mut control_query_rx,
389                    ).await;
390                }
391                Some(ipv6_packet) = dataplane_runtime.tun_outbound_rx.recv() => {
392                    let tun_budget = tun_drain_budget(LATENCY_PACKET_DRAIN_BUDGET);
393                    let mut turn = {
394                        let mut dataplane_io = dataplane_runtime.io();
395                        self.drain_dataplane_turn_with_firsts(
396                            &mut dataplane_io,
397                            crate::dataplane::DataplaneLiveTurnFirsts {
398                                tun_packet: Some(ipv6_packet),
399                                ..Default::default()
400                            },
401                            RxLoopDataplaneTurnLimits::new(0, 0, tun_budget, tun_budget),
402                        ).await
403                    };
404                    self.finish_dataplane_turn(
405                        &mut turn,
406                        &mut maintenance_state,
407                        &mut control_query_rx,
408                        0,
409                    ).await;
410                }
411                Some(identity) = dns_identity_rx.recv() => {
412                    debug!(
413                        node_addr = %identity.node_addr,
414                        "Registering identity from DNS resolution"
415                    );
416                    self.register_dns_identity(identity.node_addr, identity.pubkey);
417                }
418                Some(batch) = dataplane_runtime.endpoint_data_rx.recv() => {
419                    let mut turn = {
420                        let mut dataplane_io = dataplane_runtime.io();
421                        self.drain_dataplane_turn_with_firsts(
422                            &mut dataplane_io,
423                            crate::dataplane::DataplaneLiveTurnFirsts {
424                                endpoint_data_batch: Some(batch),
425                                ..Default::default()
426                            },
427                            RxLoopDataplaneTurnLimits::new(
428                                0,
429                                ENDPOINT_DRAIN_BUDGET,
430                                0,
431                                PACKET_DRAIN_BUDGET,
432                            ),
433                        ).await
434                    };
435                    self.finish_dataplane_turn(
436                        &mut turn,
437                        &mut maintenance_state,
438                        &mut control_query_rx,
439                        0,
440                    ).await;
441                }
442                Some((request, response_tx)) = control_command_rx.recv() => {
443                    let response = commands::dispatch(
444                        self,
445                        &request.command,
446                        request.params.as_ref(),
447                    ).await;
448                    let _ = response_tx.send(response);
449                }
450            }
451        }
452
453        info!("RX event loop stopped (channel closed)");
454        Ok(())
455    }
456
457    async fn drain_rx_loop_data_queues(
458        &mut self,
459        io: &mut RxLoopDataplaneIo<'_>,
460        budget: usize,
461    ) -> RxLoopDataDrainStats {
462        let fast_ingress =
463            Self::take_dataplane_fast_ingress_batch(io.dataplane_fast_ingress_rx, budget);
464        let packet_budget = budget.max(
465            fast_ingress
466                .as_ref()
467                .map_or(0, |fast_ingress| fast_ingress.len()),
468        );
469        let endpoint_budget = endpoint_drain_budget(packet_budget);
470        let tun_budget = tun_drain_budget(packet_budget);
471        let crypto_budget =
472            mixed_dataplane_crypto_budget(packet_budget, endpoint_budget, tun_budget);
473        let mut turn = self
474            .drain_dataplane_turn_with_firsts(
475                io,
476                crate::dataplane::DataplaneLiveTurnFirsts {
477                    fast_ingress,
478                    ..Default::default()
479                },
480                RxLoopDataplaneTurnLimits::new(
481                    packet_budget,
482                    endpoint_budget,
483                    tun_budget,
484                    crypto_budget,
485                ),
486            )
487            .await;
488        let drained_packets = Self::dataplane_packet_activity(&turn);
489        let control_drained = Box::pin(self.process_dataplane_control_ingress(&mut turn)).await;
490        RxLoopDataDrainStats::new(
491            drained_packets,
492            turn.tun_source_drained(),
493            turn.endpoint_source_drained(),
494            control_drained,
495        )
496    }
497
498    fn take_dataplane_fast_ingress_batch(
499        dataplane_fast_ingress_rx: &mut crate::dataplane::DataplaneFastIngressRx,
500        limit: usize,
501    ) -> Option<crate::dataplane::DataplaneFastIngressBatch> {
502        let fast_ingress = dataplane_fast_ingress_rx.try_recv().ok()?;
503        Some(Self::coalesce_dataplane_fast_ingress(
504            fast_ingress,
505            dataplane_fast_ingress_rx,
506            limit,
507        ))
508    }
509
510    fn coalesce_dataplane_fast_ingress(
511        mut fast_ingress: crate::dataplane::DataplaneFastIngressBatch,
512        dataplane_fast_ingress_rx: &mut crate::dataplane::DataplaneFastIngressRx,
513        limit: usize,
514    ) -> crate::dataplane::DataplaneFastIngressBatch {
515        while fast_ingress.len() < limit {
516            let Ok(next) = dataplane_fast_ingress_rx.try_recv() else {
517                break;
518            };
519            fast_ingress.absorb(next);
520        }
521        fast_ingress
522    }
523
524    async fn service_dataplane_bulk_turns(
525        &mut self,
526        io: &mut RxLoopDataplaneIo<'_>,
527        firsts: crate::dataplane::DataplaneLiveTurnFirsts,
528        maintenance_state: &mut RxLoopMaintenanceState,
529        control_query_rx: &mut Receiver<ControlMessage>,
530    ) {
531        let started = Instant::now();
532        let mut firsts = Some(firsts);
533        let mut turns = 0usize;
534
535        loop {
536            if turns > 0
537                && (turns >= RX_LOOP_BULK_SERVICE_MAX_TURNS
538                    || started.elapsed() >= RX_LOOP_BULK_SERVICE_MAX_ELAPSED
539                    || io.packet_rx.priority_ready_packets() > 0)
540            {
541                break;
542            }
543
544            let packet_budget = PACKET_DRAIN_BUDGET;
545            let mut turn_firsts = firsts.take().unwrap_or_default();
546            turn_firsts.raw_ingress_prefetch = true;
547            turn_firsts.fast_ingress = match turn_firsts.fast_ingress.take() {
548                Some(fast_ingress) => Some(Self::coalesce_dataplane_fast_ingress(
549                    fast_ingress,
550                    io.dataplane_fast_ingress_rx,
551                    packet_budget,
552                )),
553                None => Self::take_dataplane_fast_ingress_batch(
554                    io.dataplane_fast_ingress_rx,
555                    packet_budget,
556                ),
557            };
558            let packet_budget = packet_budget.max(
559                turn_firsts
560                    .fast_ingress
561                    .as_ref()
562                    .map_or(0, |fast_ingress| fast_ingress.len()),
563            );
564            let endpoint_budget = endpoint_drain_budget(packet_budget);
565            let tun_budget = tun_drain_budget(packet_budget);
566            let crypto_budget =
567                mixed_dataplane_crypto_budget(packet_budget, endpoint_budget, tun_budget);
568
569            let mut turn = self
570                .drain_dataplane_turn_with_firsts(
571                    io,
572                    turn_firsts,
573                    RxLoopDataplaneTurnLimits::new(
574                        packet_budget,
575                        endpoint_budget,
576                        tun_budget,
577                        crypto_budget,
578                    ),
579                )
580                .await;
581            let raw_drained = Self::dataplane_raw_ingress_activity(&turn);
582            let control_activity = Self::dataplane_control_activity(&turn);
583            let completions_drained = turn.summary().completions();
584            let admission_dropped =
585                turn.summary().inbound_dropped() > 0 || turn.summary().outbound_dropped() > 0;
586            let keep_servicing = !admission_dropped
587                && (raw_drained >= packet_budget
588                    || completions_drained >= crypto_budget
589                    || turn.tun_source_drained() >= tun_budget
590                    || turn.endpoint_source_drained() >= endpoint_budget);
591            let control_drained = self
592                .finish_dataplane_turn(
593                    &mut turn,
594                    maintenance_state,
595                    control_query_rx,
596                    CONTROL_QUERY_INTERLEAVE_BUDGET,
597                )
598                .await;
599            turns += 1;
600            let mut runnable_work = self.dataplane.has_runnable_work();
601
602            if control_drained == 0
603                && bulk_admission_pressure_relief_due(
604                    admission_dropped,
605                    runnable_work,
606                    turns,
607                    started.elapsed(),
608                    io.packet_rx.priority_ready_packets(),
609                )
610            {
611                let mut relief_turn = self
612                    .drain_dataplane_completion_turn(io, LATENCY_PACKET_DRAIN_BUDGET)
613                    .await;
614                let relief_control_drained = self
615                    .finish_dataplane_turn(&mut relief_turn, maintenance_state, control_query_rx, 0)
616                    .await;
617                turns += 1;
618                runnable_work = self.dataplane.has_runnable_work();
619                if relief_control_drained > 0 || !runnable_work {
620                    break;
621                }
622            }
623
624            if control_drained > 0
625                || admission_dropped
626                || (!keep_servicing && !runnable_work)
627                || (control_activity > 0 && !runnable_work)
628            {
629                break;
630            }
631        }
632    }
633
634    async fn service_dataplane_completion_turns(
635        &mut self,
636        io: &mut RxLoopDataplaneIo<'_>,
637        maintenance_state: &mut RxLoopMaintenanceState,
638        control_query_rx: &mut Receiver<ControlMessage>,
639    ) {
640        let started = Instant::now();
641        let mut turns = 0usize;
642
643        loop {
644            if turns > 0
645                && (turns >= RX_LOOP_BULK_SERVICE_MAX_TURNS
646                    || started.elapsed() >= RX_LOOP_BULK_SERVICE_MAX_ELAPSED
647                    || io.packet_rx.priority_ready_packets() > 0)
648            {
649                break;
650            }
651
652            let mut turn = self
653                .drain_dataplane_completion_turn(io, LATENCY_PACKET_DRAIN_BUDGET)
654                .await;
655            let control_drained = self
656                .finish_dataplane_turn(&mut turn, maintenance_state, control_query_rx, 0)
657                .await;
658            turns += 1;
659
660            let runnable_work = self.dataplane.has_runnable_work();
661            if control_drained > 0 || !runnable_work {
662                break;
663            }
664        }
665    }
666
667    async fn finish_dataplane_turn(
668        &mut self,
669        turn: &mut crate::dataplane::DataplaneLiveNodeTurn,
670        maintenance_state: &mut RxLoopMaintenanceState,
671        control_query_rx: &mut Receiver<ControlMessage>,
672        control_query_budget: usize,
673    ) -> usize {
674        let had_activity = turn.has_activity();
675        let control_drained = Box::pin(self.process_dataplane_control_ingress(turn))
676            .await
677            .saturating_add(Box::pin(self.drain_deferred_dataplane_control_turns()).await);
678        if control_drained > 0 && self.dataplane.has_deferred_raw_ingress() {
679            self.dataplane.readiness_notify().notify_one();
680        }
681        let query_drained = if control_query_budget > 0 {
682            self.drain_control_queries(control_query_rx, None, control_query_budget)
683                .await
684        } else {
685            0
686        };
687        if had_activity || control_drained > 0 {
688            maintenance_state.record_data_activity(Instant::now());
689        }
690        control_drained.saturating_add(query_drained)
691    }
692
693    async fn drain_control_queries(
694        &mut self,
695        control_query_rx: &mut Receiver<ControlMessage>,
696        first_message: Option<ControlMessage>,
697        budget: usize,
698    ) -> usize {
699        let mut drain = SingleLaneDrainCursor::new(first_message, budget);
700        while let Some((request, response_tx)) = drain.next(control_query_rx) {
701            let response = queries::dispatch(self, &request.command, request.params.as_ref());
702            let _ = response_tx.send(response);
703        }
704
705        drain.drained()
706    }
707
708    async fn run_rx_loop_maintenance_tick(&mut self, plan: RxLoopMaintenancePlan) -> bool {
709        if !rx_loop_fast_maintenance_within_budget(self.run_rx_loop_fast_maintenance_tick()).await {
710            crate::perf_profile::record_event(
711                crate::perf_profile::Event::RxLoopSlowMaintenanceTimeout,
712            );
713            self.mark_rx_loop_maintenance_timeout();
714            warn!(
715                timeout_ms = RX_LOOP_FAST_MAINTENANCE_TIMEOUT.as_millis() as u64,
716                data_pressure = plan.data_pressure(),
717                "RX loop liveness maintenance timed out; continuing packet processing"
718            );
719            return true;
720        }
721
722        let Some(slow_timeout) = plan.slow_timeout() else {
723            crate::perf_profile::record_event(
724                crate::perf_profile::Event::RxLoopSlowMaintenanceSkipped,
725            );
726            return false;
727        };
728
729        if tokio::time::timeout(slow_timeout, self.run_rx_loop_slow_maintenance_tick())
730            .await
731            .is_err()
732        {
733            crate::perf_profile::record_event(
734                crate::perf_profile::Event::RxLoopSlowMaintenanceTimeout,
735            );
736            self.mark_rx_loop_maintenance_timeout();
737            warn!(
738                timeout_ms = slow_timeout.as_millis() as u64,
739                data_pressure = plan.data_pressure(),
740                "RX loop slow maintenance timed out; continuing packet processing"
741            );
742            return true;
743        }
744        false
745    }
746
747    async fn run_rx_loop_fast_maintenance_tick(&mut self) {
748        self.check_timeouts();
749        let now_ms = Self::now_ms();
750        // Link/session liveness must run before slower retry/discovery work:
751        // under bulk send pressure a late heartbeat or MMP report is
752        // indistinguishable from a dead direct path on the remote peer.
753        self.check_link_heartbeats().await;
754        self.reload_peer_acl();
755        self.resend_pending_handshakes(now_ms).await;
756        self.resend_pending_rekeys(now_ms).await;
757        self.resend_pending_session_handshakes(now_ms).await;
758        self.resend_pending_session_msg3(now_ms).await;
759        self.retry_pending_session_traffic().await;
760        self.purge_idle_sessions(now_ms);
761        self.purge_learned_routes(now_ms);
762        self.check_mmp_reports().await;
763        self.check_session_mmp_reports().await;
764        self.check_rekey().await;
765        self.check_session_rekey().await;
766        self.check_pending_lookups(now_ms).await;
767        self.poll_pending_connects().await;
768        self.process_pending_retries(now_ms).await;
769        self.poll_transport_discovery().await;
770        self.sample_transport_congestion();
771    }
772
773    async fn run_rx_loop_slow_maintenance_tick(&mut self) {
774        if let Some(delay) = rx_loop_slow_maintenance_fault_delay() {
775            tokio::time::sleep(delay).await;
776        }
777
778        // Discovery and graph/stat maintenance can involve relay work or
779        // larger scans. Keep it bounded after direct-path liveness and session
780        // upkeep so a slow Nostr/LAN tick degrades discovery freshness, not
781        // packet flow.
782        self.poll_nostr_discovery().await;
783        self.poll_lan_discovery().await;
784        self.poll_local_rendezvous().await;
785        self.check_tree_state().await;
786        self.check_bloom_state().await;
787        self.compute_mesh_size();
788        self.record_stats_history();
789    }
790}
791
792async fn wait_for_optional_notify(notify: Option<&Arc<Notify>>) {
793    match notify {
794        Some(notify) => notify.notified().await,
795        None => std::future::pending().await,
796    }
797}
798
799async fn rx_loop_fast_maintenance_within_budget<F>(maintenance: F) -> bool
800where
801    F: std::future::Future<Output = ()>,
802{
803    tokio::time::timeout(RX_LOOP_FAST_MAINTENANCE_TIMEOUT, maintenance)
804        .await
805        .is_ok()
806}
807
808fn bulk_admission_pressure_relief_due(
809    admission_dropped: bool,
810    runnable_work: bool,
811    turns: usize,
812    elapsed: Duration,
813    priority_ready_packets: usize,
814) -> bool {
815    admission_dropped
816        && runnable_work
817        && turns < RX_LOOP_BULK_SERVICE_MAX_TURNS
818        && elapsed < RX_LOOP_BULK_SERVICE_MAX_ELAPSED
819        && priority_ready_packets == 0
820}