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
222 loop {
223 tokio::select! {
224 biased;
225 _ = 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 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, }
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 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 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}