From 42e94f127c6858386fb0ee5b687a9b9df83865bd Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Fri, 14 Nov 2025 22:23:18 +0530 Subject: [PATCH 1/7] feat: add a basic listener socket --- src/db/server_multithread/paxos.rs | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index 62bae88..b6a4ec5 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -1,11 +1,12 @@ -use std::net::SocketAddr; +use std::net::{SocketAddr}; use tokio::net::UdpSocket; use crate::loki_kv::control::ControlFile; pub struct ServiceManager { - udp_socket: UdpSocket, + udp_socket_send: UdpSocket, + udp_socket_recv: UdpSocket, BROADCAST_ADDRESS: SocketAddr, } @@ -14,16 +15,29 @@ impl ServiceManager { let listen_addr: SocketAddr = "0.0.0.0:8080".parse().unwrap(); let std_socket = std::net::UdpSocket::bind(listen_addr).unwrap(); + let consume_addr: SocketAddr = "0.0.0.0:8081".parse().unwrap(); + let std_consumer_socket = std::net::UdpSocket::bind(consume_addr).unwrap(); + std_socket.set_broadcast(true); ServiceManager { - udp_socket: UdpSocket::from_std(std_socket).unwrap(), + udp_socket_send: UdpSocket::from_std(std_socket).unwrap(), + udp_socket_recv: UdpSocket::from_std(std_consumer_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(); + self.udp_socket_send.send_to(msg.as_bytes(), self.BROADCAST_ADDRESS).await.unwrap(); + Ok(()) + } + + pub async fn start_consumption(self) -> Result<(), ()> { + loop{ + // TODO: Add consumption logic + // Somehitng like a go-routine treatment here? + break; + } Ok(()) } } From 40d4aad9c9aa66073ed1acfdd1c0d7d949152034 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sun, 30 Nov 2025 16:39:49 +0530 Subject: [PATCH 2/7] fix: get listen and consume addr from control file --- src/db/loki_kv/control.rs | 30 ++++++++++++++ src/db/parser/tests.rs | 65 ++++++++++++++++++++++++++++++ src/db/server_multithread/paxos.rs | 7 ++-- 3 files changed, 99 insertions(+), 3 deletions(-) create mode 100644 src/db/parser/tests.rs diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index 90ed5c7..6fa84ad 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -16,6 +16,8 @@ pub struct ControlFile { wal_directory_path: String, current_leader_value: Option, self_identifier: Option, + listen_addr: String, + consume_addr: String } impl ControlFile { @@ -26,6 +28,14 @@ impl ControlFile { self.last_wal_timeline + 1 } + pub fn get_listen_addr(&self) -> &str{ + &self.listen_addr + } + + pub fn get_consume_addr(&self) -> &str{ + &self.consume_addr + } + pub fn get_wal_directory_path(&self) -> &str { &self.wal_directory_path } @@ -64,6 +74,8 @@ impl ControlFile { last_checkpoint_id: u64, checkpoint_directory_path: String, wal_directory_path: String, + listen_addr: Option, + consume_addr: Option ) -> Result { // Create the WAL and checkpoint directories let wal_dir = Path::new(&wal_directory_path); @@ -77,6 +89,22 @@ impl ControlFile { Err(err) => return Err(err.to_string()), } + let final_listen_addr: String = match listen_addr{ + Some(addr) => addr, + None => { + info_string("no listening address provided.. defaulting to 0.0.0.0:8080".to_string()); + return "0.0.0.0:8080".to_string() + } + }; + + let final_consume_addr: String = match consume_addr{ + Some(addr) => addr, + None => { + info_string("no listening address provided.. defaulting to 0.0.0.0:8081".to_string()); + return "0.0.0.0:8081".to_string() + } + }; + let ctrl_file = ControlFile { last_wal_timeline, last_checkpoint_id, @@ -84,6 +112,8 @@ impl ControlFile { wal_directory_path, current_leader_value: None, self_identifier: Some(1 as u64), + listen_addr: Some(final_listen_addr), + consume_addr: Some(final_consume_addr) }; // Take lock on control file diff --git a/src/db/parser/tests.rs b/src/db/parser/tests.rs new file mode 100644 index 0000000..df908e9 --- /dev/null +++ b/src/db/parser/tests.rs @@ -0,0 +1,65 @@ + +#[cfg(test)] +mod tests { + use std::sync::{Arc, RwLock}; + use crate::loki_kv::loki_kv::LokiKV; + use crate::parser::parser::{parse_lokiql, QLCommands, QLValues}; + use crate::parser::executor::Executor; + + #[test] + fn test_parse_set_string() { + let query = "SET key 'value';"; + let asts = parse_lokiql(query); + assert_eq!(asts.len(), 1); + let ast = asts[0].as_ref().unwrap(); + let command_node = ast.get_left_child().unwrap(); + assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::SET))); + let key_node = command_node.get_left_child().unwrap(); + assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); + let value_node = command_node.get_right_child().unwrap(); + assert!(matches!(value_node.get_value(), QLValues::QLString(s) if s == "'value'")); + } + + #[test] + fn test_parse_set_int() { + let query = "SET key 123;"; + let asts = parse_lokiql(query); + assert_eq!(asts.len(), 1); + let ast = asts[0].as_ref().unwrap(); + let command_node = ast.get_left_child().unwrap(); + assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::SET))); + let key_node = command_node.get_left_child().unwrap(); + assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); + let value_node = command_node.get_right_child().unwrap(); + assert!(matches!(value_node.get_value(), QLValues::QLInt(123))); + } + + #[test] + fn test_parse_get() { + let query = "GET key;"; + let asts = parse_lokiql(query); + assert_eq!(asts.len(), 1); + let ast = asts[0].as_ref().unwrap(); + let command_node = ast.get_left_child().unwrap(); + assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::GET))); + let key_node = command_node.get_left_child().unwrap(); + assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); + } + + #[test] + fn test_execute_set_get() { + let db = Arc::new(RwLock::new(LokiKV::new())); + let query = "SET key 'value';"; + let asts = parse_lokiql(query); + let mut executor = Executor::new(db.clone(), asts); + executor.execute(); + + let query = "GET key;"; + let asts = parse_lokiql(query); + let mut executor = Executor::new(db.clone(), asts); + let result = executor.execute(); + + assert_eq!(result.len(), 1); + assert!(matches!(&result[0], crate::loki_kv::loki_kv::ValueObject::StringData(s) if s == "'value'")); + } +} diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index b6a4ec5..33eee51 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -2,7 +2,7 @@ use std::net::{SocketAddr}; use tokio::net::UdpSocket; -use crate::loki_kv::control::ControlFile; +use crate::loki_kv::{control::ControlFile, loki_kv::get_control_file_path}; pub struct ServiceManager { udp_socket_send: UdpSocket, @@ -12,10 +12,11 @@ pub struct ServiceManager { impl ServiceManager { pub fn new() -> Self { - let listen_addr: SocketAddr = "0.0.0.0:8080".parse().unwrap(); + let control_file = ControlFile::read_from_file_path(get_control_file_path()).unwrap(); + let listen_addr: SocketAddr = control_file.get_listen_addr().parse().unwrap(); let std_socket = std::net::UdpSocket::bind(listen_addr).unwrap(); - let consume_addr: SocketAddr = "0.0.0.0:8081".parse().unwrap(); + let consume_addr: SocketAddr = control_file.get_consume_addr().parse().unwrap(); let std_consumer_socket = std::net::UdpSocket::bind(consume_addr).unwrap(); std_socket.set_broadcast(true); From 878cc1a62c17dd841e220e8ab9630541461dadd8 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sun, 30 Nov 2025 16:40:11 +0530 Subject: [PATCH 3/7] rm: tests --- src/db/parser/tests.rs | 65 ------------------------------------------ 1 file changed, 65 deletions(-) delete mode 100644 src/db/parser/tests.rs diff --git a/src/db/parser/tests.rs b/src/db/parser/tests.rs deleted file mode 100644 index df908e9..0000000 --- a/src/db/parser/tests.rs +++ /dev/null @@ -1,65 +0,0 @@ - -#[cfg(test)] -mod tests { - use std::sync::{Arc, RwLock}; - use crate::loki_kv::loki_kv::LokiKV; - use crate::parser::parser::{parse_lokiql, QLCommands, QLValues}; - use crate::parser::executor::Executor; - - #[test] - fn test_parse_set_string() { - let query = "SET key 'value';"; - let asts = parse_lokiql(query); - assert_eq!(asts.len(), 1); - let ast = asts[0].as_ref().unwrap(); - let command_node = ast.get_left_child().unwrap(); - assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::SET))); - let key_node = command_node.get_left_child().unwrap(); - assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); - let value_node = command_node.get_right_child().unwrap(); - assert!(matches!(value_node.get_value(), QLValues::QLString(s) if s == "'value'")); - } - - #[test] - fn test_parse_set_int() { - let query = "SET key 123;"; - let asts = parse_lokiql(query); - assert_eq!(asts.len(), 1); - let ast = asts[0].as_ref().unwrap(); - let command_node = ast.get_left_child().unwrap(); - assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::SET))); - let key_node = command_node.get_left_child().unwrap(); - assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); - let value_node = command_node.get_right_child().unwrap(); - assert!(matches!(value_node.get_value(), QLValues::QLInt(123))); - } - - #[test] - fn test_parse_get() { - let query = "GET key;"; - let asts = parse_lokiql(query); - assert_eq!(asts.len(), 1); - let ast = asts[0].as_ref().unwrap(); - let command_node = ast.get_left_child().unwrap(); - assert!(matches!(command_node.get_value(), QLValues::QLCommand(QLCommands::GET))); - let key_node = command_node.get_left_child().unwrap(); - assert!(matches!(key_node.get_value(), QLValues::QLId(s) if s == "key")); - } - - #[test] - fn test_execute_set_get() { - let db = Arc::new(RwLock::new(LokiKV::new())); - let query = "SET key 'value';"; - let asts = parse_lokiql(query); - let mut executor = Executor::new(db.clone(), asts); - executor.execute(); - - let query = "GET key;"; - let asts = parse_lokiql(query); - let mut executor = Executor::new(db.clone(), asts); - let result = executor.execute(); - - assert_eq!(result.len(), 1); - assert!(matches!(&result[0], crate::loki_kv::loki_kv::ValueObject::StringData(s) if s == "'value'")); - } -} From 5713f381cc4436bac7aa88123654ace945180e51 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Wed, 3 Dec 2025 21:54:19 +0530 Subject: [PATCH 4/7] feat: spawn a message proc task on recieving message from udp socket - no testing strategy devised yet --- src/db/server_multithread/paxos.rs | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index 33eee51..873c05e 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -2,7 +2,7 @@ use std::net::{SocketAddr}; use tokio::net::UdpSocket; -use crate::loki_kv::{control::ControlFile, loki_kv::get_control_file_path}; +use crate::{loki_kv::{control::ControlFile, loki_kv::get_control_file_path}, utils::info_string}; pub struct ServiceManager { udp_socket_send: UdpSocket, @@ -37,6 +37,16 @@ impl ServiceManager { loop{ // 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); + + tokio::spawn( + async move { + // Log message + info_string(format!("Recieved the following message: {:?}", msg_bytes)); + } + ); + break; } Ok(()) From 35f3254d5ef5b1ca8ef7b2ba166c434d9fa9481c Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Fri, 5 Dec 2025 22:47:54 +0530 Subject: [PATCH 5/7] fix: use config for ports too --- src/db/loki_kv/control.rs | 22 +++++++++++++++++----- src/db/loki_kv/loki_kv.rs | 2 +- src/db/loki_kv/wal.rs | 4 ++++ src/db/main.rs | 2 +- src/db/server_multithread/server.rs | 12 +++++++++--- 5 files changed, 32 insertions(+), 10 deletions(-) diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index 6fa84ad..e3ab274 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -10,6 +10,8 @@ use crate::utils::info_string; #[derive(Serialize, Deserialize, Clone)] pub struct ControlFile { + host: String, + port: u16, last_wal_timeline: u64, last_checkpoint_id: u64, checkpoint_directory_path: String, @@ -21,6 +23,12 @@ pub struct ControlFile { } impl ControlFile { + pub fn get_hostname(&self) -> String{ + return self.host.clone(); + } + pub fn get_port(&self) -> u16{ + return self.port; + } pub fn get_next_checkpoint_id(&self) -> u64 { self.last_checkpoint_id + 1 } @@ -69,6 +77,8 @@ impl ControlFile { return false; } pub fn write( + host: String, + port: u16, path: String, last_wal_timeline: u64, last_checkpoint_id: u64, @@ -89,11 +99,11 @@ impl ControlFile { Err(err) => return Err(err.to_string()), } - let final_listen_addr: String = match listen_addr{ + let final_listen_addr: String = match listen_addr { Some(addr) => addr, None => { info_string("no listening address provided.. defaulting to 0.0.0.0:8080".to_string()); - return "0.0.0.0:8080".to_string() + "0.0.0.0:8080".to_string() } }; @@ -101,19 +111,21 @@ impl ControlFile { Some(addr) => addr, None => { info_string("no listening address provided.. defaulting to 0.0.0.0:8081".to_string()); - return "0.0.0.0:8081".to_string() + "0.0.0.0:8081".to_string() } }; let ctrl_file = ControlFile { + host, + port, last_wal_timeline, last_checkpoint_id, checkpoint_directory_path, wal_directory_path, current_leader_value: None, self_identifier: Some(1 as u64), - listen_addr: Some(final_listen_addr), - consume_addr: Some(final_consume_addr) + listen_addr: final_listen_addr, + consume_addr: final_consume_addr }; // Take lock on control file diff --git a/src/db/loki_kv/loki_kv.rs b/src/db/loki_kv/loki_kv.rs index 23d7810..e11c89d 100644 --- a/src/db/loki_kv/loki_kv.rs +++ b/src/db/loki_kv/loki_kv.rs @@ -337,7 +337,7 @@ pub fn get_data_directory() -> String { pub fn get_control_file_path() -> String { match env::var("CONTROL_FILE_PATH") { Ok(s) => s, - _ => "./control.toml".to_string(), + _ => "/home/akshat/control.toml".to_string(), } } diff --git a/src/db/loki_kv/wal.rs b/src/db/loki_kv/wal.rs index 4d9aeaf..9092d12 100644 --- a/src/db/loki_kv/wal.rs +++ b/src/db/loki_kv/wal.rs @@ -54,11 +54,15 @@ impl WALManager { pub fn new_without_toml() -> Self { let control_file = ControlFile::write( + "localhost".to_string(), + 8765 as u16, "/home/akshat/lokikv/control.toml".to_string(), 0 as u64, 0 as u64, "/home/akshat/lokikv/checkpoints".to_string(), "/home/akshat/lokikv/wal".to_string(), + Some("0.0.0.0:8080".to_string()), + Some("0.0.0.0:8081".to_string()) ) .unwrap(); let timeline = control_file.get_next_timeline_id(); diff --git a/src/db/main.rs b/src/db/main.rs index aa49068..9b7c4c2 100644 --- a/src/db/main.rs +++ b/src/db/main.rs @@ -7,6 +7,6 @@ use crate::server_multithread::server::LokiServer; #[tokio::main] async fn main() { - let serv = LokiServer::new("localhost".to_string(), 8765, 16); + let serv = LokiServer::new(16); serv.await.start_event_loop().await; } diff --git a/src/db/server_multithread/server.rs b/src/db/server_multithread/server.rs index 9e40509..09db86c 100644 --- a/src/db/server_multithread/server.rs +++ b/src/db/server_multithread/server.rs @@ -1,4 +1,5 @@ -use crate::loki_kv::loki_kv::{LokiKV, ValueObject}; +use crate::loki_kv::control::ControlFile; +use crate::loki_kv::loki_kv::{LokiKV, ValueObject, get_control_file_path}; use crate::parser::executor::Executor; use crate::parser::parser::parse_lokiql; use crate::utils::{error_string, info, info_string, warning}; @@ -23,6 +24,7 @@ pub struct LokiServer { port: u16, thread_count: usize, db_instance: Arc>, + control_file: ControlFile, } // async fn handle_connection( @@ -66,8 +68,11 @@ async fn handle_connection( } impl LokiServer { - pub async fn new(host: String, port: u16, thread_count: usize) -> Self { - let addr = format!("{}:{}", host, port); + pub async fn new(thread_count: usize) -> Self { + let control_file = ControlFile::read_from_file_path(get_control_file_path()).unwrap(); + let host: String = control_file.get_hostname(); + let port: u16 = control_file.get_port(); + let addr = format!("{}:{}", control_file.get_hostname(), control_file.get_port()); info_string(format!("Trying to start server at -> {}", addr)); let tcp_listener = TcpListener::bind(addr).await; @@ -81,6 +86,7 @@ impl LokiServer { port, thread_count, db_instance: Arc::new(RwLock::new(db_instance)), + control_file: control_file, } } Err(_) => { From 0dfa8b517a78d13238b72e752312cc4719d5de7e Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sat, 20 Dec 2025 18:51:23 +0530 Subject: [PATCH 6/7] fix: serializablity of hll --- src/db/loki_kv/data_structures/hyperloglog.rs | 3 ++- src/db/loki_kv/loki_kv.rs | 1 - 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/db/loki_kv/data_structures/hyperloglog.rs b/src/db/loki_kv/data_structures/hyperloglog.rs index 571add5..d773483 100644 --- a/src/db/loki_kv/data_structures/hyperloglog.rs +++ b/src/db/loki_kv/data_structures/hyperloglog.rs @@ -1,3 +1,4 @@ +use serde::{Deserialize, Serialize}; use std::{ collections::HashMap, fmt::Debug, @@ -9,7 +10,7 @@ use crate::loki_kv::loki_kv::ValueObject; const P_BITS: u32 = 16; const M: usize = 2_i32.pow(P_BITS) as usize; -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Serialize, Deserialize)] pub struct HLL { // leading zeros -> Count of elements streams: Vec, diff --git a/src/db/loki_kv/loki_kv.rs b/src/db/loki_kv/loki_kv.rs index e11c89d..d73379c 100644 --- a/src/db/loki_kv/loki_kv.rs +++ b/src/db/loki_kv/loki_kv.rs @@ -29,7 +29,6 @@ pub enum ValueObject { OutputString(String), BlobData(Vec), ListData(Vec), - #[serde(skip_serializing, skip_deserializing)] HLLPointer(HLL), } From 0bbac206add71dc8e809e6f61015014833eaf618 Mon Sep 17 00:00:00 2001 From: Akshat Jaimini Date: Sat, 20 Dec 2025 21:46:10 +0530 Subject: [PATCH 7/7] feat: implemented gossip protocol --- .gemini/settings.json | 8 + Cargo.lock | 219 ++++++++++++++++++++++++++-- Cargo.toml | 1 + src/db/loki_kv/control.rs | 48 ++++-- src/db/loki_kv/wal.rs | 5 +- src/db/parser/parser.rs | 31 ++-- src/db/server_multithread/paxos.rs | 75 +++++++++- src/db/server_multithread/server.rs | 71 +++++++-- 8 files changed, 406 insertions(+), 52 deletions(-) create mode 100644 .gemini/settings.json diff --git a/.gemini/settings.json b/.gemini/settings.json new file mode 100644 index 0000000..518ac41 --- /dev/null +++ b/.gemini/settings.json @@ -0,0 +1,8 @@ +{ + "mcpServers": { + "autofix": { + "command": "autofix --mcp", + "args": [] + } + } +} \ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock index 2988c48..4da5559 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -53,7 +53,7 @@ version = "1.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6d36fc52c7f6c869915e99412912f22093507da8d9e942ceaf66fe4b7c14422a" dependencies = [ - "windows-sys", + "windows-sys 0.52.0", ] [[package]] @@ -63,7 +63,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5bf74e1b6e971609db8ca7a9ce79fd5768ab6ae46441c572e46cf596f59e57f8" dependencies = [ "anstyle", - "windows-sys", + "windows-sys 0.52.0", ] [[package]] @@ -126,6 +126,12 @@ dependencies = [ "generic-array", ] +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + [[package]] name = "bytes" version = "1.8.0" @@ -228,6 +234,72 @@ dependencies = [ "typenum", ] +[[package]] +name = "darling" +version = "0.20.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc7f46116c46ff9ab3eb1597a45688b6715c6e628b5c133e288e709a29bcb4ee" +dependencies = [ + "darling_core", + "darling_macro", +] + +[[package]] +name = "darling_core" +version = "0.20.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0d00b9596d185e565c2207a0b01f8bd1a135483d02d9b7b0a54b11da8d53412e" +dependencies = [ + "fnv", + "ident_case", + "proc-macro2", + "quote", + "strsim", + "syn", +] + +[[package]] +name = "darling_macro" +version = "0.20.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc34b93ccb385b40dc71c6fceac4b2ad23662c7eeb248cf10d529b7e055b6ead" +dependencies = [ + "darling_core", + "quote", + "syn", +] + +[[package]] +name = "derive_builder" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "507dfb09ea8b7fa618fcf76e953f4f5e192547945816d5358edffe39f6f94947" +dependencies = [ + "derive_builder_macro", +] + +[[package]] +name = "derive_builder_core" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2d5bcf7b024d6835cfb3d473887cd966994907effbe9227e8c8219824d06c4e8" +dependencies = [ + "darling", + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "derive_builder_macro" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ab63b0e2bf4d5928aff72e83a7dace85d7bba5fe12dcc3c5a572d78caffd3f3c" +dependencies = [ + "derive_builder_core", + "syn", +] + [[package]] name = "digest" version = "0.10.7" @@ -250,6 +322,12 @@ version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f" +[[package]] +name = "fnv" +version = "1.0.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" + [[package]] name = "generic-array" version = "0.14.7" @@ -260,6 +338,18 @@ dependencies = [ "version_check", ] +[[package]] +name = "getset" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9cf0fc11e47561d47397154977bc219f4cf809b2974facc3ccb3b89e2436f912" +dependencies = [ + "proc-macro-error2", + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "gimli" version = "0.31.1" @@ -284,6 +374,12 @@ version = "0.3.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d231dfb89cfffdbc30e7fc41579ed6066ad03abda9e567ccafae602b97ec5024" +[[package]] +name = "ident_case" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9e0384b61958566e926dc50660321d12159025e767c18e043daf26b70104c39" + [[package]] name = "indexmap" version = "2.11.4" @@ -302,9 +398,21 @@ checksum = "7943c866cc5cd64cbc25b2e01621d07fa8eb2a1a23160ee81ce38704e97b8ecf" [[package]] name = "libc" -version = "0.2.164" +version = "0.2.178" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "433bfe06b8c75da9b2e3fbea6e5329ff87748f0b144ef75306e674c3f6f7c13f" +checksum = "37c93d8daa9d8a012fd8ab92f088405fb202ea0b6ab73ee2482ae66af4f42091" + +[[package]] +name = "local-ip-address" +version = "0.6.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a60bf300a990b2d1ebdde4228e873e8e4da40d834adbf5265f3da1457ede652" +dependencies = [ + "libc", + "neli", + "thiserror 2.0.17", + "windows-sys 0.61.2", +] [[package]] name = "lock_api" @@ -316,6 +424,12 @@ dependencies = [ "scopeguard", ] +[[package]] +name = "log" +version = "0.4.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" + [[package]] name = "lokikv" version = "0.1.0" @@ -323,6 +437,7 @@ dependencies = [ "bincode", "bit-set", "clap", + "local-ip-address", "paris", "pest", "pest_derive", @@ -357,7 +472,36 @@ dependencies = [ "hermit-abi", "libc", "wasi", - "windows-sys", + "windows-sys 0.52.0", +] + +[[package]] +name = "neli" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e23bebbf3e157c402c4d5ee113233e5e0610cc27453b2f07eefce649c7365dcc" +dependencies = [ + "bitflags", + "byteorder", + "derive_builder", + "getset", + "libc", + "log", + "neli-proc-macros", + "parking_lot", +] + +[[package]] +name = "neli-proc-macros" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "05d8d08c6e98f20a62417478ebf7be8e1425ec9acecc6f63e22da633f6b71609" +dependencies = [ + "either", + "proc-macro2", + "quote", + "serde", + "syn", ] [[package]] @@ -411,7 +555,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "879952a81a83930934cbf1786752d6dedc3b1f29e8f8fb2ad1d0a36f377cf442" dependencies = [ "memchr", - "thiserror", + "thiserror 1.0.69", "ucd-trie", ] @@ -455,6 +599,28 @@ version = "0.2.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "915a1e146535de9163f3987b8944ed8cf49a18bb0056bcebcdcece385cece4ff" +[[package]] +name = "proc-macro-error-attr2" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "96de42df36bb9bba5542fe9f1a054b8cc87e172759a1868aa05c1f3acc89dfc5" +dependencies = [ + "proc-macro2", + "quote", +] + +[[package]] +name = "proc-macro-error2" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "11ec05c52be0a07b08061f7dd003e7d7092e0472bc731b4af7bb1ef876109802" +dependencies = [ + "proc-macro-error-attr2", + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "proc-macro2" version = "1.0.86" @@ -592,7 +758,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ce305eb0b4296696835b71df73eb912e0f1ffd2556a501fcede6e0c50349191c" dependencies = [ "libc", - "windows-sys", + "windows-sys 0.52.0", ] [[package]] @@ -618,7 +784,16 @@ version = "1.0.69" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" dependencies = [ - "thiserror-impl", + "thiserror-impl 1.0.69", +] + +[[package]] +name = "thiserror" +version = "2.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f63587ca0f12b72a0600bcba1d40081f830876000bb46dd2337a3051618f4fc8" +dependencies = [ + "thiserror-impl 2.0.17", ] [[package]] @@ -632,6 +807,17 @@ dependencies = [ "syn", ] +[[package]] +name = "thiserror-impl" +version = "2.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3ff15c8ecd7de3849db632e14d18d2571fa09dfc5ed93479bc4485c7a517c913" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "tokio" version = "1.41.1" @@ -647,7 +833,7 @@ dependencies = [ "signal-hook-registry", "socket2", "tokio-macros", - "windows-sys", + "windows-sys 0.52.0", ] [[package]] @@ -736,6 +922,12 @@ version = "0.11.0+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + [[package]] name = "windows-sys" version = "0.52.0" @@ -745,6 +937,15 @@ dependencies = [ "windows-targets", ] +[[package]] +name = "windows-sys" +version = "0.61.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" +dependencies = [ + "windows-link", +] + [[package]] name = "windows-targets" version = "0.52.6" diff --git a/Cargo.toml b/Cargo.toml index b98e12b..8e9a8be 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -24,3 +24,4 @@ pest_derive = "2.6" bit-set = "0.8.0" paris = { version = "1.5", features = ["timestamps"] } toml = "0.9.7" +local-ip-address = "0.6.8" diff --git a/src/db/loki_kv/control.rs b/src/db/loki_kv/control.rs index e3ab274..efecec4 100644 --- a/src/db/loki_kv/control.rs +++ b/src/db/loki_kv/control.rs @@ -18,8 +18,11 @@ pub struct ControlFile { wal_directory_path: String, current_leader_value: Option, self_identifier: Option, - listen_addr: String, - consume_addr: String + send_addr: String, + consume_addr: String, + checkpoint_timer_interval: Option, + paxos_timer_interval: Option, + gossip_timeout: Option } impl ControlFile { @@ -36,8 +39,8 @@ impl ControlFile { self.last_wal_timeline + 1 } - pub fn get_listen_addr(&self) -> &str{ - &self.listen_addr + pub fn get_send_addr(&self) -> &str{ + &self.send_addr } pub fn get_consume_addr(&self) -> &str{ @@ -60,6 +63,27 @@ impl ControlFile { self.self_identifier.clone() } + pub fn get_checkpoint_timer_interval(&self) -> u64 { + match self.checkpoint_timer_interval{ + Some(val) => return val, + None => return 1 + } + } + + pub fn get_paxos_timer_interval(&self) -> u64 { + match self.paxos_timer_interval{ + Some(val) => return val, + None => return 5 + } + } + + pub fn get_gossip_timeout(&self) -> u64{ + match self.gossip_timeout{ + Some(val) => return val, + None => return 300 + } + } + pub fn set_current_leader_identifier(&mut self, current_leader_value: u64) { self.current_leader_value = Some(current_leader_value.clone()); } @@ -84,8 +108,11 @@ impl ControlFile { last_checkpoint_id: u64, checkpoint_directory_path: String, wal_directory_path: String, - listen_addr: Option, - consume_addr: Option + send_addr: Option, + consume_addr: Option, + checkpoint_timer_interval: Option, + paxos_timer_interval: Option, + gossip_timeout: Option ) -> Result { // Create the WAL and checkpoint directories let wal_dir = Path::new(&wal_directory_path); @@ -99,7 +126,7 @@ impl ControlFile { Err(err) => return Err(err.to_string()), } - let final_listen_addr: String = match listen_addr { + let final_send_addr: String = match send_addr { Some(addr) => addr, None => { info_string("no listening address provided.. defaulting to 0.0.0.0:8080".to_string()); @@ -124,8 +151,11 @@ impl ControlFile { wal_directory_path, current_leader_value: None, self_identifier: Some(1 as u64), - listen_addr: final_listen_addr, - consume_addr: final_consume_addr + send_addr: final_send_addr, + consume_addr: final_consume_addr, + checkpoint_timer_interval, + paxos_timer_interval, + gossip_timeout }; // Take lock on control file diff --git a/src/db/loki_kv/wal.rs b/src/db/loki_kv/wal.rs index 9092d12..4714708 100644 --- a/src/db/loki_kv/wal.rs +++ b/src/db/loki_kv/wal.rs @@ -62,7 +62,10 @@ impl WALManager { "/home/akshat/lokikv/checkpoints".to_string(), "/home/akshat/lokikv/wal".to_string(), Some("0.0.0.0:8080".to_string()), - Some("0.0.0.0:8081".to_string()) + Some("0.0.0.0:8081".to_string()), + None, + None, + None ) .unwrap(); let timeline = control_file.get_next_timeline_id(); diff --git a/src/db/parser/parser.rs b/src/db/parser/parser.rs index e45f1a2..1a0c0f0 100644 --- a/src/db/parser/parser.rs +++ b/src/db/parser/parser.rs @@ -1,3 +1,4 @@ +use crate::utils::error; use std::ops::Deref; use pest::iterators::Pair; @@ -99,21 +100,27 @@ impl AST { } pub fn parse_lokiql(ql: &str) -> Vec> { - let result = LokiQLParser::parse(Rule::LOKIQL_FILE, ql).unwrap(); - - let mut asts: Vec> = vec![]; - for pair in result { - match pair.as_rule() { - // Parse Each command - Rule::COMMAND => { - let ast = parse_vals(pair, None); - asts.push(ast); + let result = LokiQLParser::parse(Rule::LOKIQL_FILE, ql); + match result { + Ok(pairs) => { + let mut asts: Vec> = vec![]; + for pair in pairs { + match pair.as_rule() { + // Parse Each command + Rule::COMMAND => { + let ast = parse_vals(pair, None); + asts.push(ast); + } + _ => {} + } } - _ => {} + asts + } + Err(e) => { + error(&format!("Error parsing LokiQL: {}", e.to_string())); + vec![] } } - - return asts; } pub fn parse_individual_item_asql(pair: Pair) -> QLValues { diff --git a/src/db/server_multithread/paxos.rs b/src/db/server_multithread/paxos.rs index 873c05e..de870a8 100644 --- a/src/db/server_multithread/paxos.rs +++ b/src/db/server_multithread/paxos.rs @@ -1,29 +1,37 @@ -use std::net::{SocketAddr}; +use std::{collections::{HashMap, HashSet}, io::Split, net::SocketAddr}; +use local_ip_address::local_ip; -use tokio::net::UdpSocket; +use tokio::{net::UdpSocket, time::timeout}; +use tokio::time::Duration; -use crate::{loki_kv::{control::ControlFile, loki_kv::get_control_file_path}, utils::info_string}; +use crate::{loki_kv::{control::ControlFile, loki_kv::get_control_file_path}, utils::{info_string, warning, warning_string}}; + +// ---------------------------- SERVICE MANAGER ------------------------------------------- pub struct ServiceManager { udp_socket_send: UdpSocket, udp_socket_recv: UdpSocket, + node_directory: HashSet<(String, String)>, // Hashset of node_id, address(ip + port) BROADCAST_ADDRESS: SocketAddr, } 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_listen_addr().parse().unwrap(); + let listen_addr: SocketAddr = control_file.get_send_addr().parse().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 mut node_directory: HashSet<(String, String)> = HashSet::new(); + std_socket.set_broadcast(true); 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(), } } @@ -33,7 +41,7 @@ impl ServiceManager { Ok(()) } - pub async fn start_consumption(self) -> Result<(), ()> { + pub async fn start_consumption(&self) -> Result<(), ()> { loop{ // TODO: Add consumption logic // Somehitng like a go-routine treatment here? @@ -51,6 +59,19 @@ impl ServiceManager { } Ok(()) } + + pub fn update_node_directory(&mut self, node_id: String, node_addr: String) { + self.node_directory.insert((node_id, node_addr)); + } +} + +// ---------------- PAXOS NODE ---------------------------- + +fn get_ip_addr(addr: String) -> String{ + let my_local_ip = local_ip().unwrap(); + let mut tks = addr.split(":"); + let ip = format!("{}:{}", my_local_ip.to_string(), tks.nth(1).unwrap().to_string()); + return ip; } pub struct PaxosNode { @@ -59,7 +80,11 @@ pub struct PaxosNode { } impl PaxosNode { - pub fn propose(&self) { + pub fn new_node() -> Self{ + let control_file = ControlFile::read_from_file_path(get_control_file_path()).unwrap(); + PaxosNode { ctrl_file: control_file, service_manager: ServiceManager::new()} + } + pub async fn propose(&self) { let value = self.ctrl_file.get_self_identifier(); // Broadcast to network @@ -74,4 +99,42 @@ impl PaxosNode { pub fn accept() {} pub fn learn() {} + + // Gossip for node discovery + pub async fn gossip(&self) -> Result<(), String> { + let node_id = self.ctrl_file.get_self_identifier().unwrap(); + let data = format!("{}~{}", node_id.to_string(), get_ip_addr(self.ctrl_file.get_consume_addr().to_string())); + info_string(format!("Sending -> {}", data)); + let result = self.service_manager.broadcast_message(data.as_str()).await; + return result; + } + + pub async fn gossip_consume(&mut self) { + let MAX_GOSSIP_CONSUMPTION = 10; + for i in 0..MAX_GOSSIP_CONSUMPTION{ + info_string(format!("{} gossip trial", i)); + // Somehitng like a go-routine treatment here? + let mut msg_bytes: Vec = vec![]; + match timeout(Duration::from_secs(self.ctrl_file.get_gossip_timeout()), self.service_manager.udp_socket_recv.recv_from(&mut msg_bytes)).await{ + Ok(Ok((_, _))) => { + let data = String::from_utf8(msg_bytes).unwrap(); + let mut tokens; + if data.contains("~"){ + tokens = data.as_str().split("~"); + }else{ + let msg = format!("Token does not contain any ~ {:?}, skipping..", data); + warning(msg.as_str()); + continue; + } + + // Log message + info_string(format!("Recieved the following message: {:?}", tokens)); + self.service_manager.update_node_directory(tokens.nth(0).unwrap().to_string(), tokens.nth(1).unwrap().to_string()); + }, + Ok(Err(e)) => panic!("{}", e), + Err(e) => panic!("{}", e) + }; + + } + } } diff --git a/src/db/server_multithread/server.rs b/src/db/server_multithread/server.rs index 09db86c..60042dc 100644 --- a/src/db/server_multithread/server.rs +++ b/src/db/server_multithread/server.rs @@ -2,6 +2,7 @@ use crate::loki_kv::control::ControlFile; use crate::loki_kv::loki_kv::{LokiKV, ValueObject, get_control_file_path}; use crate::parser::executor::Executor; use crate::parser::parser::parse_lokiql; +use crate::server_multithread::paxos::PaxosNode; use crate::utils::{error_string, info, info_string, warning}; use std::env; use std::time::{Duration, Instant}; @@ -49,21 +50,34 @@ async fn handle_connection( // .map_err(|e| format!("Invalid UTF-8 data: {}", e)) // .unwrap(); + let mut resp_str = String::new(); + let asts = parse_lokiql(&request_line); - let mut ast_exector = Executor::new(db_instance.clone(), asts); - let responses = ast_exector.execute(); + if asts.len() == 0{ + // Query was wrong.. lets tell it to the user + resp_str += "Invalid command.. Pls try again\n"; + }else{ + let mut ast_exector = Executor::new(db_instance.clone(), asts); + let responses = ast_exector.execute(); - let mut resp_str = String::new(); - // Improve output result - for response in responses.iter() { - if let val = response { - resp_str += &format!("{:?}\n", val); - }; + // Improve output result + for response in responses.iter() { + if let val = response { + resp_str += &format!("{:?}\n", val); + }; + } } resp_str += "\n"; - let _ = wr.write_all(resp_str.as_bytes()).await; - let _ = wr.flush().await; + + wr.write_all(resp_str.as_bytes()).await + .map_err(|e| format!("Failed to write response: {}", e))?; + + wr.flush().await + .map_err(|e| format!("Failed to flush writer: {}", e))?; + + info_string(format!("Sent response: {} bytes", resp_str)); + } } @@ -96,12 +110,16 @@ impl LokiServer { } pub async fn start_event_loop(&mut self) { - let checkpoint_itr: u64 = match env::var("CHECKPOINT_INTERVAL") { - Ok(n) => n.parse().unwrap(), - _ => 120 as u64, - }; + let checkpoint_itr: u64 = self.control_file.get_checkpoint_timer_interval(); + let paxos_itr: u64 = self.control_file.get_paxos_timer_interval(); + + let mut checkpoint_timer = interval(Duration::from_secs(checkpoint_itr*60)); + let mut paxos_gossip_broadcast_timer = interval(Duration::from_secs(paxos_itr*30)); + // let mut paxos_gossip_consumer_timer = interval(Duration::from_secs(paxos_itr*60)); - let mut checkpoint_timer = interval(Duration::from_secs(checkpoint_itr)); + let mut paxos_node = PaxosNode::new_node(); + + let mut should_broadcast = true; loop { select! { @@ -121,6 +139,7 @@ impl LokiServer { } } + _ = checkpoint_timer.tick() => { info("Checkpointing..."); let ins = self.db_instance.clone(); @@ -129,6 +148,28 @@ impl LokiServer { db.checkpoint(); }); } + + _ = paxos_gossip_broadcast_timer.tick() => { + // Send out self information 10 times + if should_broadcast{ + for _ in 0..10{ + info("Gossiping self information with strangers!"); + let res = paxos_node.gossip().await; + match res{ + Ok(_) => info("Successfully shared gossip information."), + Err(err) => info_string(format!("Failed to share gossip information {}", err)) + }; + } + sleep(Duration::from_secs(10)).await; + }else{ + // Consume gossip data + info("Consuming gossip from strangers..(for now)"); + paxos_node.gossip_consume().await; + } + + should_broadcast = !should_broadcast; + + } } } }