From a06f6bbd7543e16c3174f84f501dc48b58627d11 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sun, 1 Feb 2026 22:32:40 +0530 Subject: [PATCH 1/2] use soc 2 to enable reuse addr --- src/db/server_multithread/paxos.rs | 19 ++++++++++++++++--- 1 file changed, 16 insertions(+), 3 deletions(-) diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index 5cae146..b10e1b6 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -1,4 +1,5 @@ use local_ip_address::local_ip; +use socket2::{Domain, SockAddr, Socket, Type}; use std::{ collections::{HashMap, HashSet}, io::Split, @@ -26,14 +27,26 @@ impl ServiceManager { pub fn new() -> Self { let control_file = ControlFile::read_from_file_path(get_control_file_path()).unwrap(); let listen_addr: SocketAddr = control_file.get_send_addr().parse().unwrap(); + let soc2_listen_socket = Socket::new(Domain::IPV4, Type::DGRAM, None).unwrap(); + soc2_listen_socket.set_broadcast(true).unwrap(); + soc2_listen_socket + .bind(&SockAddr::from(listen_addr)) + .unwrap(); let std_socket = std::net::UdpSocket::bind(listen_addr).unwrap(); let consume_addr: SocketAddr = control_file.get_consume_addr().parse().unwrap(); - let std_consumer_socket = std::net::UdpSocket::bind(consume_addr).unwrap(); + let soc2_raw_consumer_socket = match Socket::new(Domain::IPV4, Type::DGRAM, None) { + Ok(socket) => socket, + Err(e) => panic!("Failed to create socket: {}", e), + }; + soc2_raw_consumer_socket.set_reuse_address(true).unwrap(); + soc2_raw_consumer_socket + .bind(&SockAddr::from(consume_addr)) + .unwrap(); - let mut node_directory: HashSet<(String, String)> = HashSet::new(); + let std_consumer_socket: std::net::UdpSocket = soc2_raw_consumer_socket.into(); - std_socket.set_broadcast(true); + let node_directory: HashSet<(String, String)> = HashSet::new(); ServiceManager { udp_socket_send: UdpSocket::from_std(std_socket).unwrap(), From 1ca5c0a5d7d40cf6f56089802a7c4bbe515fc139 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sun, 1 Feb 2026 23:07:03 +0530 Subject: [PATCH 2/2] fix: allow nodes to run on same port --- src/db/server_multithread/paxos.rs | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index b10e1b6..d417147 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -32,7 +32,7 @@ impl ServiceManager { soc2_listen_socket .bind(&SockAddr::from(listen_addr)) .unwrap(); - let std_socket = std::net::UdpSocket::bind(listen_addr).unwrap(); + let std_socket: std::net::UdpSocket = soc2_listen_socket.into(); let consume_addr: SocketAddr = control_file.get_consume_addr().parse().unwrap(); let soc2_raw_consumer_socket = match Socket::new(Domain::IPV4, Type::DGRAM, None) { @@ -48,11 +48,13 @@ impl ServiceManager { let node_directory: HashSet<(String, String)> = HashSet::new(); + let broadcast_ip_string = format!("255.255.255.255:{}", consume_addr.port()); + ServiceManager { udp_socket_send: UdpSocket::from_std(std_socket).unwrap(), udp_socket_recv: UdpSocket::from_std(std_consumer_socket).unwrap(), node_directory: node_directory, - BROADCAST_ADDRESS: "255.255.255.255:8080".parse().unwrap(), + BROADCAST_ADDRESS: broadcast_ip_string.parse().unwrap(), } } @@ -69,7 +71,10 @@ impl ServiceManager { // TODO: Add consumption logic // Somehitng like a go-routine treatment here? let mut msg_bytes: Vec = vec![]; - self.udp_socket_recv.recv_from(&mut msg_bytes); + self.udp_socket_recv + .recv_from(&mut msg_bytes) + .await + .unwrap(); tokio::spawn(async move { // Log message