2021-12-03 09:00:31 -08:00
|
|
|
use {
|
2023-09-11 09:57:10 -07:00
|
|
|
crate::repair::{quic_endpoint::RemoteRequest, serve_repair::ServeRepair},
|
|
|
|
crossbeam_channel::{unbounded, Receiver, Sender},
|
2021-12-03 09:00:31 -08:00
|
|
|
solana_ledger::blockstore::Blockstore,
|
2023-09-11 09:57:10 -07:00
|
|
|
solana_perf::{packet::PacketBatch, recycler::Recycler},
|
2022-05-05 11:56:18 -07:00
|
|
|
solana_streamer::{
|
|
|
|
socket::SocketAddrSpace,
|
|
|
|
streamer::{self, StreamerReceiveStats},
|
|
|
|
},
|
2021-12-03 09:00:31 -08:00
|
|
|
std::{
|
|
|
|
net::UdpSocket,
|
2022-08-01 11:22:54 -07:00
|
|
|
sync::{atomic::AtomicBool, Arc},
|
2023-09-11 09:57:10 -07:00
|
|
|
thread::{self, Builder, JoinHandle},
|
2023-03-31 08:42:49 -07:00
|
|
|
time::Duration,
|
2021-12-03 09:00:31 -08:00
|
|
|
},
|
|
|
|
};
|
2020-01-31 14:23:51 -08:00
|
|
|
|
|
|
|
pub struct ServeRepairService {
|
|
|
|
thread_hdls: Vec<JoinHandle<()>>,
|
|
|
|
}
|
|
|
|
|
|
|
|
impl ServeRepairService {
|
|
|
|
pub fn new(
|
2022-08-01 11:22:54 -07:00
|
|
|
serve_repair: ServeRepair,
|
2023-09-11 09:57:10 -07:00
|
|
|
remote_request_sender: Sender<RemoteRequest>,
|
|
|
|
remote_request_receiver: Receiver<RemoteRequest>,
|
2022-08-01 11:46:45 -07:00
|
|
|
blockstore: Arc<Blockstore>,
|
2020-01-31 14:23:51 -08:00
|
|
|
serve_repair_socket: UdpSocket,
|
2021-07-23 08:25:03 -07:00
|
|
|
socket_addr_space: SocketAddrSpace,
|
2021-12-17 15:21:05 -08:00
|
|
|
stats_reporter_sender: Sender<Box<dyn FnOnce() + Send>>,
|
2022-08-01 11:22:54 -07:00
|
|
|
exit: Arc<AtomicBool>,
|
2020-01-31 14:23:51 -08:00
|
|
|
) -> Self {
|
2022-01-11 02:44:46 -08:00
|
|
|
let (request_sender, request_receiver) = unbounded();
|
2020-01-31 14:23:51 -08:00
|
|
|
let serve_repair_socket = Arc::new(serve_repair_socket);
|
|
|
|
trace!(
|
|
|
|
"ServeRepairService: id: {}, listening on: {:?}",
|
2022-08-01 11:22:54 -07:00
|
|
|
&serve_repair.my_id(),
|
2020-01-31 14:23:51 -08:00
|
|
|
serve_repair_socket.local_addr().unwrap()
|
|
|
|
);
|
|
|
|
let t_receiver = streamer::receiver(
|
|
|
|
serve_repair_socket.clone(),
|
2022-05-05 11:56:18 -07:00
|
|
|
exit.clone(),
|
2020-01-31 14:23:51 -08:00
|
|
|
request_sender,
|
2021-04-07 08:15:38 -07:00
|
|
|
Recycler::default(),
|
2022-05-05 11:56:18 -07:00
|
|
|
Arc::new(StreamerReceiveStats::new("serve_repair_receiver")),
|
2023-03-31 08:42:49 -07:00
|
|
|
Duration::from_millis(1), // coalesce
|
2023-09-11 09:57:10 -07:00
|
|
|
false, // use_pinned_memory
|
|
|
|
None, // in_vote_only_mode
|
2020-01-31 14:23:51 -08:00
|
|
|
);
|
2023-09-11 09:57:10 -07:00
|
|
|
let t_packet_adapter = Builder::new()
|
|
|
|
.name(String::from("solServRAdapt"))
|
|
|
|
.spawn(|| adapt_repair_requests_packets(request_receiver, remote_request_sender))
|
|
|
|
.unwrap();
|
2022-01-11 02:44:46 -08:00
|
|
|
let (response_sender, response_receiver) = unbounded();
|
2021-07-23 08:25:03 -07:00
|
|
|
let t_responder = streamer::responder(
|
2022-08-17 08:40:23 -07:00
|
|
|
"Repair",
|
2021-07-23 08:25:03 -07:00
|
|
|
serve_repair_socket,
|
|
|
|
response_receiver,
|
|
|
|
socket_addr_space,
|
2021-12-17 15:21:05 -08:00
|
|
|
Some(stats_reporter_sender),
|
2021-07-23 08:25:03 -07:00
|
|
|
);
|
2023-09-11 09:57:10 -07:00
|
|
|
let t_listen =
|
|
|
|
serve_repair.listen(blockstore, remote_request_receiver, response_sender, exit);
|
2020-01-31 14:23:51 -08:00
|
|
|
|
2023-09-11 09:57:10 -07:00
|
|
|
let thread_hdls = vec![t_receiver, t_packet_adapter, t_responder, t_listen];
|
2020-01-31 14:23:51 -08:00
|
|
|
Self { thread_hdls }
|
|
|
|
}
|
|
|
|
|
2023-08-28 14:34:09 -07:00
|
|
|
pub(crate) fn join(self) -> thread::Result<()> {
|
|
|
|
self.thread_hdls.into_iter().try_for_each(JoinHandle::join)
|
2020-01-31 14:23:51 -08:00
|
|
|
}
|
|
|
|
}
|
2023-09-11 09:57:10 -07:00
|
|
|
|
|
|
|
// Adapts incoming UDP repair requests into RemoteRequest struct.
|
|
|
|
pub(crate) fn adapt_repair_requests_packets(
|
|
|
|
packets_receiver: Receiver<PacketBatch>,
|
|
|
|
remote_request_sender: Sender<RemoteRequest>,
|
|
|
|
) {
|
|
|
|
for packets in packets_receiver {
|
|
|
|
for packet in &packets {
|
|
|
|
let Some(bytes) = packet.data(..).map(Vec::from) else {
|
|
|
|
continue;
|
|
|
|
};
|
|
|
|
let request = RemoteRequest {
|
|
|
|
remote_pubkey: None,
|
|
|
|
remote_address: packet.meta().socket_addr(),
|
|
|
|
bytes,
|
|
|
|
response_sender: None,
|
|
|
|
};
|
|
|
|
if remote_request_sender.send(request).is_err() {
|
|
|
|
return; // The receiver end of the channel is disconnected.
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|