Contractless/src/rpc/server/start_rpc.rs

106 lines
3.6 KiB
Rust
Raw Normal View History

2026-07-25 14:50:39 +00:00
use crate::log::{error, warn};
2026-05-24 17:56:57 +00:00
use crate::records::memory::response_channels::Command;
use crate::rpc::server::handshake::handle_handshake;
use crate::sled::Db;
use crate::wallets::structures::Wallet;
2026-05-24 17:56:57 +00:00
use crate::Arc;
use crate::Mutex;
use crate::SocketAddr;
use crate::TcpListener;
2026-07-25 14:50:39 +00:00
use socket2::{Domain, Protocol, Socket, Type};
2026-05-24 17:56:57 +00:00
// wait incomming connections
pub async fn start_rpc(
db: &Db,
server_address: String,
wallet: Arc<Wallet>,
2026-05-24 17:56:57 +00:00
map: Arc<Mutex<Command>>,
) {
// Parse once at startup so the accept loop can work with a concrete
// socket address instead of reparsing the settings string each time.
let server_socket: SocketAddr = server_address
.parse()
.expect("Failed to parse server address");
let db_clone = db.clone();
// The listener runs in the background while startup continues with
// the rest of node initialization.
tokio::spawn(async move {
rpc_server(server_socket, &db_clone, wallet, map).await;
2026-05-24 17:56:57 +00:00
});
}
// generate a connection when it comes in
async fn rpc_server(
server_socket: SocketAddr,
db: &Db,
wallet: Arc<Wallet>,
2026-05-24 17:56:57 +00:00
map: Arc<Mutex<Command>>,
) {
// Bind failure means this node cannot accept RPC traffic, so log the
// reason and leave the background task instead of panicking.
let listener = match TcpListener::bind(&server_socket).await {
Ok(listener) => listener,
Err(e) => {
error!("Failed to bind to socket: {e:?}");
return;
}
};
2026-07-25 14:50:39 +00:00
// IPv4 remains the node transport and advertised identity. A second
// IPv6-only listener accepts RPC clients on the same configured port;
// the handshake rejects any IPv6 socket that advertises a miner port.
if server_socket.is_ipv4() {
let ipv6_listener = (|| -> std::io::Result<TcpListener> {
let socket = Socket::new(Domain::IPV6, Type::STREAM, Some(Protocol::TCP))?;
socket.set_only_v6(true)?;
let address = SocketAddr::new("::".parse().unwrap(), server_socket.port());
socket.bind(&address.into())?;
socket.listen(1024)?;
let listener: std::net::TcpListener = socket.into();
listener.set_nonblocking(true)?;
TcpListener::from_std(listener)
})();
match ipv6_listener {
Ok(ipv6_listener) => {
let ipv6_db = db.clone();
let ipv6_wallet = wallet.clone();
let ipv6_map = map.clone();
tokio::spawn(async move {
accept_connections(ipv6_listener, ipv6_db, ipv6_wallet, ipv6_map).await;
});
}
Err(err) => {
warn!("IPv6 RPC client listener unavailable: {err}");
}
}
}
accept_connections(listener, db.clone(), wallet, map).await;
}
async fn accept_connections(
listener: TcpListener,
db: Db,
wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>,
) {
2026-05-24 17:56:57 +00:00
loop {
match listener.accept().await {
Ok((stream, _)) => {
// Every accepted socket gets its own handshake task so
// slow peers do not block the listener from accepting.
let stream = Arc::new(Mutex::new(stream));
let db_clone = db.clone();
let wallet_clone = wallet.clone();
2026-05-24 17:56:57 +00:00
let map_clone = map.clone();
tokio::spawn(async move {
handle_handshake(stream, db_clone, wallet_clone, map_clone).await;
2026-05-24 17:56:57 +00:00
});
}
Err(e) => {
error!("Error accepting connection: {e:?}");
}
}
}
}