From 94b67b644cd135262f35626ceaff66a94790c032 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Wed, 12 Nov 2025 22:12:01 +0530 Subject: [PATCH 1/5] feat: init paxos implementaion --- src/db/server_multithread/mod.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/db/server_multithread/mod.rs b/src/db/server_multithread/mod.rs index 74f47ad..d497856 100644 --- a/src/db/server_multithread/mod.rs +++ b/src/db/server_multithread/mod.rs @@ -1 +1,2 @@ pub mod server; +pub mod paxos; From 989c7c50529daddc312d677f7d13463fafbee2a7 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Wed, 12 Nov 2025 22:13:07 +0530 Subject: [PATCH 2/5] feat: add paxos interface --- src/db/server_multithread/paxos.rs | 7 +++++++ 1 file changed, 7 insertions(+) create mode 100644 src/db/server_multithread/paxos.rs diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs new file mode 100644 index 0000000..c584570 --- /dev/null +++ b/src/db/server_multithread/paxos.rs @@ -0,0 +1,7 @@ +pub struct PaxosNode{} + +impl PaxosNode{ + pub fn propose(){} + pub fn accept(){} + pub fn learn() {} +} From 26869cfa48889db2d5cf0388ebfb0027e081f70f Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Wed, 12 Nov 2025 22:15:40 +0530 Subject: [PATCH 3/5] feat: add current leader value to control file --- src/db/loki_kv/control.rs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index f5d5c5f..832935a 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -14,6 +14,7 @@ pub struct ControlFile { last_checkpoint_id: u64, checkpoint_directory_path: String, wal_directory_path: String, + current_leader_value: Option, } impl ControlFile { @@ -29,6 +30,9 @@ impl ControlFile { pub fn get_checkpoint_directory_path(&self) -> &str { &self.checkpoint_directory_path } + pub fn get_current_leader_identifier(&self) -> Option{ + self.current_leader_value.clone() + } pub fn write( path: String, last_wal_timeline: u64, From f32d677457ff5e96c8db112bfad9d94dd04ee172 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Wed, 12 Nov 2025 22:23:09 +0530 Subject: [PATCH 4/5] feat: update control file to store leader values --- src/db/loki_kv/control.rs | 25 +++++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index 832935a..dc3a820 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -15,6 +15,7 @@ pub struct ControlFile { checkpoint_directory_path: String, wal_directory_path: String, current_leader_value: Option, + self_identifier: Option, } impl ControlFile { @@ -24,15 +25,39 @@ impl ControlFile { pub fn get_next_timeline_id(&self) -> u64 { self.last_wal_timeline + 1 } + pub fn get_wal_directory_path(&self) -> &str { &self.wal_directory_path } + pub fn get_checkpoint_directory_path(&self) -> &str { &self.checkpoint_directory_path } + pub fn get_current_leader_identifier(&self) -> Option{ self.current_leader_value.clone() } + + pub fn get_self_identifier(&self) -> Option{ + self.self_identifier.clone() + } + + pub fn set_current_leader_identifier(&mut self, current_leader_value: u64){ + self.current_leader_value = Some(current_leader_value.clone()); + } + + pub fn set_self_identifier(&mut self, id: u64) { + self.self_identifier = Some(id.clone()); + } + + pub fn is_leader(&self) -> bool{ + if self.self_identifier.is_some() && self.current_leader_value.is_some(){ + if self.self_identifier.unwrap() == self.current_leader_value.unwrap(){ + return true; + } + } + return false; + } pub fn write( path: String, last_wal_timeline: u64, From 6c105559bcd52d091adb52d9822e86332ae32dad Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Thu, 13 Nov 2025 23:31:21 +0530 Subject: [PATCH 5/5] feat: implement propose functionality --- src/db/loki_kv/control.rs | 14 ++++---- src/db/server_multithread/mod.rs | 2 +- src/db/server_multithread/paxos.rs | 53 +++++++++++++++++++++++++++--- 3 files changed, 58 insertions(+), 11 deletions(-) diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index dc3a820..90ed5c7 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -34,15 +34,15 @@ impl ControlFile { &self.checkpoint_directory_path } - pub fn get_current_leader_identifier(&self) -> Option{ + pub fn get_current_leader_identifier(&self) -> Option { self.current_leader_value.clone() } - pub fn get_self_identifier(&self) -> Option{ + pub fn get_self_identifier(&self) -> Option { self.self_identifier.clone() } - pub fn set_current_leader_identifier(&mut self, current_leader_value: u64){ + pub fn set_current_leader_identifier(&mut self, current_leader_value: u64) { self.current_leader_value = Some(current_leader_value.clone()); } @@ -50,9 +50,9 @@ impl ControlFile { self.self_identifier = Some(id.clone()); } - pub fn is_leader(&self) -> bool{ - if self.self_identifier.is_some() && self.current_leader_value.is_some(){ - if self.self_identifier.unwrap() == self.current_leader_value.unwrap(){ + pub fn is_leader(&self) -> bool { + if self.self_identifier.is_some() && self.current_leader_value.is_some() { + if self.self_identifier.unwrap() == self.current_leader_value.unwrap() { return true; } } @@ -82,6 +82,8 @@ impl ControlFile { last_checkpoint_id, checkpoint_directory_path, wal_directory_path, + current_leader_value: None, + self_identifier: Some(1 as u64), }; // Take lock on control file diff --git a/src/db/server_multithread/mod.rs b/src/db/server_multithread/mod.rs index d497856..9556874 100644 --- a/src/db/server_multithread/mod.rs +++ b/src/db/server_multithread/mod.rs @@ -1,2 +1,2 @@ -pub mod server; pub mod paxos; +pub mod server; diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index c584570..62bae88 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -1,7 +1,52 @@ -pub struct PaxosNode{} +use std::net::SocketAddr; -impl PaxosNode{ - pub fn propose(){} - pub fn accept(){} +use tokio::net::UdpSocket; + +use crate::loki_kv::control::ControlFile; + +pub struct ServiceManager { + udp_socket: UdpSocket, + BROADCAST_ADDRESS: SocketAddr, +} + +impl ServiceManager { + pub fn new() -> Self { + let listen_addr: SocketAddr = "0.0.0.0:8080".parse().unwrap(); + let std_socket = std::net::UdpSocket::bind(listen_addr).unwrap(); + + std_socket.set_broadcast(true); + + ServiceManager { + udp_socket: UdpSocket::from_std(std_socket).unwrap(), + BROADCAST_ADDRESS: "255.255.255.255:8080".parse().unwrap(), + } + } + + pub async fn broadcast_message(&self, msg: &str) -> Result<(), String> { + self.udp_socket.send_to(msg.as_bytes(), self.BROADCAST_ADDRESS).await.unwrap(); + Ok(()) + } +} + +pub struct PaxosNode { + ctrl_file: ControlFile, + service_manager: ServiceManager, +} + +impl PaxosNode { + pub fn propose(&self) { + let value = self.ctrl_file.get_self_identifier(); + + // Broadcast to network + match value{ + Some(val) => { + let msg = format!("PROPOSE {}", val); + self.service_manager.broadcast_message(msg.as_str()); + }, + None => panic!("No value to broadcast!") + } + } + + pub fn accept() {} pub fn learn() {} }