@@ -34,12 +34,14 @@ use std::thread;
3434use std:: time:: { Duration , Instant } ;
3535use tokio:: sync:: broadcast:: Receiver ;
3636
37+ const REVERSE_DNS_LOOKUP_THREADS : usize = 5 ;
38+
3739/// The calling thread enters a loop in which it waits for network packets
3840#[ allow( clippy:: too_many_lines, clippy:: too_many_arguments) ]
3941pub fn parse_packets (
4042 cap_id : usize ,
4143 mut cs : CaptureSource ,
42- mmdb_readers : MmdbReaders ,
44+ mmdb_readers : & MmdbReaders ,
4345 ip_blacklist : & IpBlacklist ,
4446 capture_context : CaptureContext ,
4547 filters : Filters ,
@@ -59,15 +61,21 @@ pub fn parse_packets(
5961
6062 let mut info_traffic_msg = InfoTraffic :: default ( ) ;
6163
62- let ( lookup_request_tx, lookup_request_rx) = std :: sync :: mpsc :: channel ( ) ;
64+ let ( lookup_request_tx, lookup_request_rx) = async_channel :: unbounded ( ) ;
6365 let ( lookup_result_tx, lookup_result_rx) = std:: sync:: mpsc:: channel ( ) ;
6466 let mut resolutions_state = AddressesResolutionState :: new ( lookup_request_tx, lookup_result_rx) ;
65- let _ = thread:: Builder :: new ( )
66- . name ( "thread_reverse_dns_lookups" . to_string ( ) )
67- . spawn ( move || {
68- reverse_dns_lookups ( & lookup_request_rx, & lookup_result_tx, & mmdb_readers) ;
69- } )
70- . log_err ( location ! ( ) ) ;
67+ // a pool of threads shares the request queue, so one slow blocking lookup doesn't stall the others
68+ for i in 0 ..REVERSE_DNS_LOOKUP_THREADS {
69+ let lookup_request_rx = lookup_request_rx. clone ( ) ;
70+ let lookup_result_tx = lookup_result_tx. clone ( ) ;
71+ let mmdb_readers = mmdb_readers. clone ( ) ;
72+ let _ = thread:: Builder :: new ( )
73+ . name ( format ! ( "thread_reverse_dns_lookups_{i}" ) )
74+ . spawn ( move || {
75+ reverse_dns_lookups ( & lookup_request_rx, & lookup_result_tx, & mmdb_readers) ;
76+ } )
77+ . log_err ( location ! ( ) ) ;
78+ }
7179
7280 // instant of the first parsed packet plus multiples of 1 second (only used in live captures)
7381 let mut first_packet_ticks = None ;
@@ -213,8 +221,8 @@ pub fn parse_packets(
213221 DataInfo :: new_with_first_packet ( exchanged_bytes, traffic_direction) ,
214222 ) ;
215223
216- // send the rDNS lookup request to the corresponding thread
217- let _ = resolutions_state. lookup_request_tx . send ( (
224+ // send the rDNS lookup request to the thread pool
225+ let _ = resolutions_state. lookup_request_tx . try_send ( (
218226 key,
219227 traffic_direction,
220228 cs. get_addresses ( ) . clone ( ) ,
@@ -360,15 +368,12 @@ fn from_linux_sll(packet: &[u8], is_v1: bool) -> Option<LaxPacketHeaders<'_>> {
360368}
361369
362370fn reverse_dns_lookups (
363- lookup_request_rx : & std:: sync:: mpsc:: Receiver < (
364- AddressPortPair ,
365- TrafficDirection ,
366- Vec < Address > ,
367- ) > ,
371+ lookup_request_rx : & async_channel:: Receiver < ( AddressPortPair , TrafficDirection , Vec < Address > ) > ,
368372 lookup_result_tx : & std:: sync:: mpsc:: Sender < HostMessage > ,
369373 mmdb_readers : & MmdbReaders ,
370374) {
371- while let Ok ( ( key, traffic_direction, interface_addresses) ) = lookup_request_rx. recv ( ) {
375+ while let Ok ( ( key, traffic_direction, interface_addresses) ) = lookup_request_rx. recv_blocking ( )
376+ {
372377 let address_to_lookup = get_address_to_lookup ( & key, traffic_direction) ;
373378
374379 // perform rDNS lookup
@@ -418,7 +423,7 @@ fn reverse_dns_lookups(
418423}
419424
420425pub struct AddressesResolutionState {
421- lookup_request_tx : std :: sync :: mpsc :: Sender < ( AddressPortPair , TrafficDirection , Vec < Address > ) > ,
426+ lookup_request_tx : async_channel :: Sender < ( AddressPortPair , TrafficDirection , Vec < Address > ) > ,
422427 lookup_result_rx : std:: sync:: mpsc:: Receiver < HostMessage > ,
423428 /// Map of the addresses waiting for a rDNS resolution; used to NOT send multiple rDNS for the same address
424429 addresses_waiting_resolution : HashMap < IpAddr , DataInfo > ,
@@ -428,11 +433,7 @@ pub struct AddressesResolutionState {
428433
429434impl AddressesResolutionState {
430435 fn new (
431- lookup_request_tx : std:: sync:: mpsc:: Sender < (
432- AddressPortPair ,
433- TrafficDirection ,
434- Vec < Address > ,
435- ) > ,
436+ lookup_request_tx : async_channel:: Sender < ( AddressPortPair , TrafficDirection , Vec < Address > ) > ,
436437 lookup_result_rx : std:: sync:: mpsc:: Receiver < HostMessage > ,
437438 ) -> Self {
438439 Self {
0 commit comments