1use 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 pub async fn run_rx_loop(&mut self) -> Result<(), NodeError> {
113 let packet_rx = self.packet_rx.take().ok_or(NodeError::NotStarted)?;
114
115 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 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 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 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(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 crate::perf_profile::maybe_spawn_reporter();
217 tick.tick().await;
221 let mut nostr_event_turn_not_before = Instant::now();
222
223 loop {
224 tokio::select! {
225 biased;
226 _ = 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 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 _ = 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, }
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 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 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}