From 901cff543f0ede80df5d93f57120b956e37a9779 Mon Sep 17 00:00:00 2001 From: moyukhb Date: Mon, 10 Mar 2025 23:02:10 +0530 Subject: [PATCH] Added serialization and deserialization logic --- .idea/.gitignore | 3 + .idea/LokiKV.iml | 9 +++ .idea/misc.xml | 6 ++ .idea/modules.xml | 8 ++ .idea/vcs.xml | 6 ++ Cargo.lock | 25 ++++++ Cargo.toml | 7 +- src/cli/main.rs | 10 ++- src/db/server_multithread/deserializer.rs | 11 +++ src/db/server_multithread/mod.rs | 2 + src/db/server_multithread/serializer.rs | 6 ++ src/db/server_multithread/server.rs | 98 ++++++++++++----------- 12 files changed, 142 insertions(+), 49 deletions(-) create mode 100644 .idea/.gitignore create mode 100644 .idea/LokiKV.iml create mode 100644 .idea/misc.xml create mode 100644 .idea/modules.xml create mode 100644 .idea/vcs.xml create mode 100644 src/db/server_multithread/deserializer.rs create mode 100644 src/db/server_multithread/serializer.rs diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..26d3352 --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,3 @@ +# Default ignored files +/shelf/ +/workspace.xml diff --git a/.idea/LokiKV.iml b/.idea/LokiKV.iml new file mode 100644 index 0000000..d6ebd48 --- /dev/null +++ b/.idea/LokiKV.iml @@ -0,0 +1,9 @@ + + + + + + + + + \ No newline at end of file diff --git a/.idea/misc.xml b/.idea/misc.xml new file mode 100644 index 0000000..6f29fee --- /dev/null +++ b/.idea/misc.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/.idea/modules.xml b/.idea/modules.xml new file mode 100644 index 0000000..0631590 --- /dev/null +++ b/.idea/modules.xml @@ -0,0 +1,8 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..35eb1dd --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock index 104b482..1769034 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -278,6 +278,12 @@ version = "1.70.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7943c866cc5cd64cbc25b2e01621d07fa8eb2a1a23160ee81ce38704e97b8ecf" +[[package]] +name = "itoa" +version = "1.0.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" + [[package]] name = "libc" version = "0.2.164" @@ -305,6 +311,7 @@ dependencies = [ "pest_derive", "rayon", "serde", + "serde_json", "shlex", "tokio", ] @@ -478,6 +485,12 @@ version = "0.1.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "719b953e2095829ee67db738b3bfa9fa368c94900df327b3f07fe6e794d2fe1f" +[[package]] +name = "ryu" +version = "1.0.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "28d3b2b1366ec20994f1fd18c3c594f05c5dd4bc44d8bb0c1c632c8d6829481f" + [[package]] name = "scopeguard" version = "1.2.0" @@ -504,6 +517,18 @@ dependencies = [ "syn", ] +[[package]] +name = "serde_json" +version = "1.0.140" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "20068b6e96dc6c9bd23e01df8827e6c7e1f2fddd43c21810382803c136b99373" +dependencies = [ + "itoa", + "memchr", + "ryu", + "serde", +] + [[package]] name = "sha2" version = "0.10.8" diff --git a/Cargo.toml b/Cargo.toml index 6a0a903..0ee2f20 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,11 +14,12 @@ path = "src/cli/main.rs" [dependencies] bincode = "1.3.3" +bit-set = "0.8.0" clap = { version = "4.5.17", features = ["derive"] } +pest = "2.6" +pest_derive = "2.6" rayon = "1.10.0" serde = { version = "1.0", features = ["derive"] } +serde_json = "1.0" shlex = "1.3.0" tokio = { version = "1.41.1", features = ["full"] } -pest = "2.6" -pest_derive = "2.6" -bit-set = "0.8.0" diff --git a/src/cli/main.rs b/src/cli/main.rs index 962223c..0c87404 100644 --- a/src/cli/main.rs +++ b/src/cli/main.rs @@ -67,10 +67,18 @@ fn main() { } // println!("Writing to stream: {}", buf); - if let Err(e) = writer.write_all(buf.as_bytes()) { + let json_command = format!(r#"{{"query": "{}"}}"#, buf.trim()); + + if let Err(e) = writer.write_all(json_command.as_bytes()) { eprintln!("Failed to send command: {}", e); break; } + + // Ensure a newline is sent for proper deserialization + if let Err(e) = writer.write_all(b"\n") { + eprintln!("Failed to send newline: {}", e); + break; + } // println!("Written to stream!"); // println!("Checking response...."); diff --git a/src/db/server_multithread/deserializer.rs b/src/db/server_multithread/deserializer.rs new file mode 100644 index 0000000..5c95501 --- /dev/null +++ b/src/db/server_multithread/deserializer.rs @@ -0,0 +1,11 @@ +use serde::Deserialize; +use serde_json::{self, Result}; + +#[derive(Debug, Deserialize)] +pub struct Request { + pub query: String, +} + +pub fn deserialize(input: &str) -> Result { + serde_json::from_str(input) +} \ No newline at end of file diff --git a/src/db/server_multithread/mod.rs b/src/db/server_multithread/mod.rs index 74f47ad..8c1eb33 100644 --- a/src/db/server_multithread/mod.rs +++ b/src/db/server_multithread/mod.rs @@ -1 +1,3 @@ pub mod server; +pub mod serializer; +pub mod deserializer; \ No newline at end of file diff --git a/src/db/server_multithread/serializer.rs b/src/db/server_multithread/serializer.rs new file mode 100644 index 0000000..275b437 --- /dev/null +++ b/src/db/server_multithread/serializer.rs @@ -0,0 +1,6 @@ +use serde::Serialize; +use serde_json; + +pub fn serialize(data: &T) -> Result { + serde_json::to_string(data).map_err(|e| format!("Serialization error: {}", e)) +} \ No newline at end of file diff --git a/src/db/server_multithread/server.rs b/src/db/server_multithread/server.rs index bde6b6f..e3c7b49 100644 --- a/src/db/server_multithread/server.rs +++ b/src/db/server_multithread/server.rs @@ -1,15 +1,25 @@ -use crate::loki_kv::loki_kv::{LokiKV, ValueObject}; +use crate::loki_kv::loki_kv::LokiKV; use crate::parser::executor::Executor; use crate::parser::parser::parse_lokiql; +use crate::server_multithread::serializer::serialize; +use crate::server_multithread::deserializer::deserialize; +use serde::{Deserialize, Serialize}; +use serde_json::Value; use std::{ - ops::{Deref, DerefMut}, sync::{Arc, RwLock}, }; -use tokio::io::{AsyncBufReadExt, BufReader}; -use tokio::{ - io::{self, AsyncReadExt, AsyncWriteExt}, - net::{TcpListener, TcpStream}, -}; +use tokio::io::{AsyncBufReadExt, BufReader, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream}; + +#[derive(Debug, Serialize, Deserialize)] +pub struct Request { + query: String, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct Response { + result: Vec, +} // Server Logic pub struct LokiServer { @@ -19,7 +29,7 @@ pub struct LokiServer { thread_count: usize, db_instance: Arc>, } -// + async fn handle_connection( stream: TcpStream, db_instance: Arc>, @@ -38,29 +48,35 @@ async fn handle_connection( } let request_line = buf.trim().to_string(); - // let request_line = String::from_utf8(buf[..n].to_vec()) - // .map_err(|e| format!("Invalid UTF-8 data: {}", e)) - // .unwrap(); - println!("Got {:?}", request_line); - let asts = parse_lokiql(&request_line); - let mut ast_exector = Executor::new(db_instance.clone(), asts); - let responses = ast_exector.execute(); + // Fix deserialization + match deserialize(&request_line) { + Ok(query) => { + println!("Executing query: {:?}", query); - let mut resp_str = String::new(); - // Improve output result - for response in responses.iter() { - if let val = response { - resp_str += &format!("{:?}\n", val); - }; - } + let asts = parse_lokiql(&query.query); + let mut ast_exector = Executor::new(db_instance.clone(), asts); + let responses = ast_exector.execute(); - resp_str += "\n"; - println!("RESPONSE: {}", resp_str); - let _ = wr.write_all(resp_str.as_bytes()).await; - let _ = wr.flush().await; - // println!("Wrote bytes {}", resp_str); + let mut resp_str = String::new(); + for response in responses.iter() { + if let val = response { + resp_str += &format!("{:?}\n", val); + }; + } + + resp_str += "\n"; + println!("RESPONSE: {}", resp_str); + let _ = wr.write_all(resp_str.as_bytes()).await; + let _ = wr.flush().await; + } + Err(e) => { + eprintln!("Deserialization error: {}", e); + let _ = wr.write_all(b"Deserialization error\n").await; + let _ = wr.flush().await; + } + } } } @@ -68,29 +84,21 @@ impl LokiServer { pub async fn new(host: String, port: u16, thread_count: usize) -> Self { let addr = format!("{}:{}", host, port); println!("Trying to start server at -> {}", addr); - let tcp_listener = TcpListener::bind(addr).await; + let tcp_listener = TcpListener::bind(addr).await.expect("Unable to create server"); - match tcp_listener { - Ok(tcp_list) => { - println!("Started Sevrer at {}:{}", host, port); - let db_instance = LokiKV::new(); - LokiServer { - tcp_listener: tcp_list, - host, - port, - thread_count, - db_instance: Arc::new(RwLock::new(db_instance)), - } - } - Err(_) => { - panic!("Unable to Create new Server at {}:{}", host, port); - } + println!("Started Server at {}:{}", host, port); + let db_instance = LokiKV::new(); + LokiServer { + tcp_listener, + host, + port, + thread_count, + db_instance: Arc::new(RwLock::new(db_instance)), } } pub async fn start_event_loop(&mut self) { loop { - // println!("HII!!"); match self.tcp_listener.accept().await { Ok((socket, _)) => { let db = self.db_instance.clone(); @@ -101,7 +109,7 @@ impl LokiServer { } }); } - _ => panic!("error accepting connection"), + Err(e) => eprintln!("Error accepting connection: {}", e), }; } }