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),
};
}
}