diff --git a/engine.rs b/engine.rs new file mode 100644 index 0000000..89d20e8 --- /dev/null +++ b/engine.rs @@ -0,0 +1,31 @@ +// engine.rs +use crate::iterator::ScanIterator; + +// ... existing code ... + +pub fn scan(&self) -> ScanIterator { + ScanIterator::new(&self.memtable, &self.sstables) +} + +pub fn keys(&self) -> Vec { + const MAX_SCAN_LIMIT: usize = 1000; // or configurable + self.scan() + .take(MAX_SCAN_LIMIT) + .map(|(key, _)| key) + .collect() +} + +pub fn count(&self) -> usize { + let memtable_count = self.memtable.len(); + let sstable_count: usize = self.sstables.iter().map(|s| s.record_count()).sum(); + memtable_count + sstable_count +} + +// search() also needs updating to use iterator pattern +pub fn search(&self, query: &str) -> Vec<(String, (Vec, u128, bool))> { + self.scan() + .take(1000) // also cap search results + .filter(|(_, (_, _, deleted))| !*deleted) + .filter(|(_, (value, _, _))| String::from_utf8_lossy(value).contains(query)) + .collect() +} diff --git a/iterator.rs b/iterator.rs new file mode 100644 index 0000000..c852912 --- /dev/null +++ b/iterator.rs @@ -0,0 +1,98 @@ +// iterator.rs +use crate::memtable::MemTable; +use crate::table::SSTable; +use std::cmp::Ordering; +use std::collections::BinaryHeap; + +pub struct ScanIterator { + heap: BinaryHeap, +} + +struct ScanItem { + key: String, + value: (Vec, u128, bool), + source: Source, +} + +enum Source { + MemTable, + SSTable(usize), // index of SSTable +} + +impl PartialEq for ScanItem { + fn eq(&self, other: &Self) -> bool { + self.key == other.key + } +} + +impl Eq for ScanItem {} + +impl PartialOrd for ScanItem { + fn partial_cmp(&self, other: &Self) -> Option { + // Reverse for min-heap (BinaryHeap is max-heap by default) + other.key.partial_cmp(&self.key) + } +} + +impl Ord for ScanItem { + fn cmp(&self, other: &Self) -> Ordering { + other.key.cmp(&self.key) + } +} + +impl ScanIterator { + pub fn new(memtable: &MemTable, sstables: &[SSTable]) -> Self { + let mut heap = BinaryHeap::new(); + + // Add all memtable entries + for (key, value) in memtable.iter() { + heap.push(ScanItem { + key: key.clone(), + value: value.clone(), + source: Source::MemTable, + }); + } + + // Add first entry from each SSTable + for (idx, sstable) in sstables.iter().enumerate() { + if let Some((key, value)) = sstable.first_key_value() { + heap.push(ScanItem { + key: key.clone(), + value: value.clone(), + source: Source::SSTable(idx), + }); + } + } + + ScanIterator { heap } + } +} + +impl Iterator for ScanIterator { + type Item = (String, (Vec, u128, bool)); + + fn next(&mut self) -> Option { + let item = self.heap.pop()?; + let key = item.key; + let value = item.value; + + // Advance the iterator from which this item came + match item.source { + Source::MemTable => { + // MemTable is already fully in heap, nothing more to add + } + Source::SSTable(idx) => { + // Get next entry from this SSTable + if let Some((next_key, next_value)) = SSTable::next_at_index(idx) { + self.heap.push(ScanItem { + key: next_key, + value: next_value, + source: Source::SSTable(idx), + }); + } + } + } + + Some((key, value)) + } +} diff --git a/memtable.rs b/memtable.rs new file mode 100644 index 0000000..e3b6603 --- /dev/null +++ b/memtable.rs @@ -0,0 +1,10 @@ +// memtable.rs +impl MemTable { + pub fn iter(&self) -> impl Iterator, u128, bool))> { + self.map.iter() + } + + pub fn len(&self) -> usize { + self.map.len() + } +} diff --git a/src/api/mod.rs b/src/api/mod.rs index 27aaa6b..93a9551 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -1,24 +1,49 @@ -use crate::core::engine::{DEFAULT_SCAN_LIMIT, MAX_SCAN_LIMIT}; +pub mod auth; +pub mod config; + +pub use self::config::ServerConfig; use crate::LsmEngine; +use actix_web::{get, web, App, HttpResponse, HttpServer, Responder}; +use serde_json::json; -pub struct ServerConfig; -impl ServerConfig { - pub fn from_file(_path: &str) -> crate::infra::error::Result { - Ok(Self) - } - pub fn from_env() -> Self { - Self +/// Handler for `GET /keys`. +/// Returns a JSON object containing an array of all keys (bounded by `MAX_SCAN_LIMIT`). +#[get("/keys")] +async fn get_keys(engine: web::Data) -> impl Responder { + // `LsmEngine::keys` applies the safety bound (MAX_SCAN_LIMIT). + match engine.keys() { + Ok(keys) => HttpResponse::Ok() + .content_type("application/json") + .json(json!({ "keys": keys })), + Err(e) => { + eprintln!("Failed to fetch keys: {:?}", e); + HttpResponse::InternalServerError() + .content_type("application/json") + .json(json!({ "error": "internal server error" })) + } } } -pub async fn start_server( - _engine: LsmEngine, - _config: ServerConfig, -) -> crate::infra::error::Result<()> { - Ok(()) +/// Register API routes. +pub fn configure(cfg: &mut web::ServiceConfig) { + cfg.service(get_keys); } -fn _check_api() { - let _ = DEFAULT_SCAN_LIMIT; - let _ = MAX_SCAN_LIMIT; +/// Start the REST API server. +pub async fn start_server(engine: LsmEngine, config: ServerConfig) -> std::io::Result<()> { + let host = config.host.clone(); + let port = config.port; + + println!("πŸš€ Starting server at http://{}:{}", host, port); + + let engine_data = web::Data::new(engine); + + HttpServer::new(move || { + App::new() + .app_data(engine_data.clone()) + .configure(configure) + }) + .bind((host, port))? + .run() + .await } diff --git a/src/bin/server.rs b/src/bin/server.rs index 99485f1..53d9eca 100644 --- a/src/bin/server.rs +++ b/src/bin/server.rs @@ -59,7 +59,9 @@ async fn main() -> std::io::Result<()> { .sparse_index_interval(sparse_index_interval) .bloom_false_positive_rate(bloom_false_positive_rate) .build() - .map_err(|e| io::Error::new(io::ErrorKind::InvalidInput, e.to_string()))?; + .map_err(|e: apexstore::LsmError| { + io::Error::new(io::ErrorKind::InvalidInput, e.to_string()) + })?; // Print LSM configuration println!("πŸ“‹ LSM Engine Configuration:"); @@ -92,5 +94,5 @@ async fn main() -> std::io::Result<()> { apexstore::api::start_server(engine, server_config) .await - .map_err(|e| io::Error::other(e.to_string())) + .map_err(|e: io::Error| e) } diff --git a/src/bin/tui.rs b/src/bin/tui.rs index 1cd0500..3d19f0d 100644 --- a/src/bin/tui.rs +++ b/src/bin/tui.rs @@ -8,7 +8,7 @@ //! SEARCH [--prefix] | SCAN | ALL | KEYS | COUNT //! STATS [ALL] | BATCH | BATCH SET | DEMO | CLEAR | HELP -use apexstore::{core::engine::LsmStats, infra::config::LsmConfig, LsmEngine}; +use apexstore::{LsmConfig, LsmEngine, LsmError, LsmStats}; use chrono::Local; use crossterm::{ event::{ @@ -34,6 +34,8 @@ use std::{ use tui_input::backend::crossterm::EventHandler; use tui_input::Input; +type ScanResult = Result, Vec)>, LsmError>; + // ─── Palette ────────────────────────────────────────────────────────────────── const C_ORANGE: Color = Color::Rgb(255, 110, 30); const C_AMBER: Color = Color::Rgb(255, 185, 0); @@ -170,11 +172,12 @@ impl App { } match self.engine.get(parts[1]) { Ok(Some(v)) => { + let value: Vec = v; self.log_push( format!( "\u{2713} '{}' = '{}'", parts[1], - String::from_utf8_lossy(&v) + String::from_utf8_lossy(&value) ), C_OK, ); @@ -274,46 +277,54 @@ impl App { } // ALL ────────────────────────────────────────────────────────────── - "ALL" => match self.engine.scan() { - Ok(rows) if rows.is_empty() => self.log_push("\u{26a0} Database is empty", C_WARN), - Ok(rows) => { - self.log_push(format!("\u{2713} {} record(s):", rows.len()), C_OK); - for (k, v) in rows.iter().take(30) { - self.log_push( - format!( - " {} = {}", - String::from_utf8_lossy(k), - String::from_utf8_lossy(v) - ), - C_TEXT, - ); + "ALL" => { + let res: ScanResult = self.engine.scan(); + match res { + Ok(rows) if rows.is_empty() => { + self.log_push("\u{26a0} Database is empty", C_WARN) } - if rows.len() > 30 { - self.log_push(format!(" ... and {} more", rows.len() - 30), C_DIM); + Ok(rows) => { + self.log_push(format!("\u{2713} {} record(s):", rows.len()), C_OK); + for (k, v) in rows.iter().take(30) { + self.log_push( + format!( + " {} = {}", + String::from_utf8_lossy(k), + String::from_utf8_lossy(v) + ), + C_TEXT, + ); + } + if rows.len() > 30 { + self.log_push(format!(" ... and {} more", rows.len() - 30), C_DIM); + } + self.incr_ops(); } - self.incr_ops(); + Err(e) => self.log_push(format!("\u{274c} {}", e), C_ERR), } - Err(e) => self.log_push(format!("\u{274c} {}", e), C_ERR), - }, + } // KEYS ───────────────────────────────────────────────────────────── - "KEYS" => match self.engine.keys() { - Ok(keys) if keys.is_empty() => self.log_push("\u{26a0} No keys found", C_WARN), - Ok(keys) => { - self.log_push(format!("\u{2713} {} key(s):", keys.len()), C_OK); - for (i, k) in keys.iter().enumerate().take(30) { - self.log_push( - format!(" {}. {}", i + 1, String::from_utf8_lossy(k)), - C_TEXT, - ); - } - if keys.len() > 30 { - self.log_push(format!(" ... and {} more", keys.len() - 30), C_DIM); + "KEYS" => { + let res: Result>, LsmError> = self.engine.keys(); + match res { + Ok(keys) if keys.is_empty() => self.log_push("\u{26a0} No keys found", C_WARN), + Ok(keys) => { + self.log_push(format!("\u{2713} {} key(s):", keys.len()), C_OK); + for (i, k) in keys.iter().enumerate().take(30) { + self.log_push( + format!(" {}. {}", i + 1, String::from_utf8_lossy(k)), + C_TEXT, + ); + } + if keys.len() > 30 { + self.log_push(format!(" ... and {} more", keys.len() - 30), C_DIM); + } + self.incr_ops(); } - self.incr_ops(); + Err(e) => self.log_push(format!("\u{274c} {}", e), C_ERR), } - Err(e) => self.log_push(format!("\u{274c} {}", e), C_ERR), - }, + } // COUNT ──────────────────────────────────────────────────────────── "COUNT" => match self.engine.count() { @@ -585,9 +596,9 @@ fn main() -> io::Result<()> { .dir_path(PathBuf::from("./.lsm_data")) .memtable_max_size(64 * 1024) // 64 KB .build() - .map_err(|e| io::Error::other(e.to_string()))?; + .map_err(|e: LsmError| io::Error::other(e.to_string()))?; - let engine = LsmEngine::new(config).map_err(|e| io::Error::other(e.to_string()))?; + let engine = LsmEngine::new(config).map_err(|e: LsmError| io::Error::other(e.to_string()))?; let mut terminal = setup()?; let mut app = App::new(engine); diff --git a/src/core/engine/mod.rs b/src/core/engine/mod.rs index 0a8fb76..0321caf 100644 --- a/src/core/engine/mod.rs +++ b/src/core/engine/mod.rs @@ -9,6 +9,8 @@ use std::collections::HashMap; use self::compaction::Compaction; use self::manifest::Manifest; use self::version_set::VersionSet; +use crate::core::iterators::{MergeIterator, StorageIterator}; +use crate::core::key::KeySlice; pub const DEFAULT_SCAN_LIMIT: usize = 128; pub const MAX_SCAN_LIMIT: usize = 1024; @@ -109,9 +111,40 @@ impl MemTable { self.size -= old.len(); } } +} + +struct InternalMemTableIterator<'a> { + inner: std::collections::btree_map::Iter<'a, Vec, Vec>, + current: Option<(&'a Vec, &'a Vec)>, +} + +impl<'a> InternalMemTableIterator<'a> { + fn new(data: &'a std::collections::BTreeMap, Vec>) -> Self { + let mut inner = data.iter(); + let current = inner.next(); + Self { inner, current } + } +} + +impl<'a> StorageIterator for InternalMemTableIterator<'a> { + type KeyType = KeySlice<'a>; - fn iter(&self) -> impl Iterator, Vec)> + '_ { - self.data.iter().map(|(k, v)| (k.clone(), v.clone())) + fn next(&mut self) { + self.current = self.inner.next(); + } + fn key(&self) -> Self::KeyType { + KeySlice::new(self.current.unwrap().0.as_slice()) + } + fn value(&self) -> &[u8] { + self.current.unwrap().1.as_slice() + } + fn is_valid(&self) -> bool { + self.current.is_some() + } + fn seek(&mut self, _key: &[u8]) { + while self.is_valid() && self.key().as_ref() < _key { + self.next(); + } } } @@ -236,38 +269,52 @@ impl Engine { upper: Option<&[u8]>, limit: Option, ) -> crate::infra::error::Result, Vec)>> { - let mut seen_keys = std::collections::HashSet::new(); - let mut results = Vec::new(); + let mut iters: Vec> + '_>> = Vec::new(); - // 1. Memtables primeiro (mais recentes) + // 1. Memtables (newer first) if let Some(memtables) = self.memtables.get(cf) { for mem in memtables.iter().rev() { - for (k, v) in mem.iter() { - // aplicar filtros lower/upper - if lower.is_none_or(|lb| k.as_slice() >= lb) - && upper.is_none_or(|ub| k.as_slice() < ub) - && seen_keys.insert(k.clone()) - { - results.push((k.clone(), v.clone())); - } - if let Some(l) = limit { - if results.len() >= l { - return Ok(results); - } - } - } + iters.push(Box::new(InternalMemTableIterator::new(&mem.data))); } } - // 2. SSTables do version_set (mais antigas) - let sst_results = self.version_set.scan(cf, lower, upper, limit); - for (k, v) in sst_results { - if seen_keys.insert(k.clone()) { - results.push((k, v)); + // 2. SSTables (from VersionSet) + for sst_iter in self.version_set.table_iters(cf) { + iters.push(Box::new(sst_iter)); + } + + let mut merge_iter = MergeIterator::new(iters); + + // Se houver search bound inferior, seek + if let Some(lb) = lower { + // Nota: Nosso MergeIterator.seek ainda nΓ£o estΓ‘ implementado, mas para scans bΓ‘sicos + // podemos apenas skipar. No futuro, seek otimizado seria preferΓ­vel. + while merge_iter.is_valid() && merge_iter.key().as_slice() < lb { + merge_iter.next(); } - if limit.is_some_and(|l| results.len() >= l) { - break; + } + + let mut results = Vec::new(); + while merge_iter.is_valid() { + let key = merge_iter.key(); + let key_slice: &[u8] = key.as_ref(); + + // Check upper bound + if let Some(ub) = upper { + if key_slice >= ub { + break; + } + } + + // Apenas adicionar se nΓ£o for tombstone (valor vazio em algumas impls, mas aqui vamos assumir todos vΓ‘lidos por enquanto) + results.push((key_slice.to_vec(), merge_iter.value().to_vec())); + + if let Some(l) = limit { + if results.len() >= l { + break; + } } + merge_iter.next(); } Ok(results) @@ -314,16 +361,38 @@ impl Engine { } pub fn keys(&self) -> crate::infra::error::Result>> { - Ok(self.memtables.get("default").map_or(Vec::new(), |m| { - m.iter().flat_map(|mt| mt.data.keys().cloned()).collect() - })) + let mut iters: Vec> + '_>> = Vec::new(); + + if let Some(memtables) = self.memtables.get("default") { + for mem in memtables.iter().rev() { + iters.push(Box::new(InternalMemTableIterator::new(&mem.data))); + } + } + + for sst_iter in self.version_set.table_iters("default") { + iters.push(Box::new(sst_iter)); + } + + let mut merge_iter = MergeIterator::new(iters); + let mut results = Vec::new(); + + while merge_iter.is_valid() && results.len() < MAX_SCAN_LIMIT { + results.push(merge_iter.key()); + merge_iter.next(); + } + + Ok(results) } pub fn count(&self) -> crate::infra::error::Result { - Ok(self + let mem_count: usize = self .memtables .get("default") - .map_or(0, |m| m.iter().map(|mt| mt.data.len()).sum())) + .map_or(0, |m| m.iter().map(|mt| mt.data.len()).sum()); + + let sst_count = self.version_set.record_count("default"); + + Ok(mem_count + sst_count) } pub fn stats(&self) -> crate::infra::error::Result { diff --git a/src/core/engine/version_set.rs b/src/core/engine/version_set.rs index 7060c9a..9d7646d 100644 --- a/src/core/engine/version_set.rs +++ b/src/core/engine/version_set.rs @@ -63,6 +63,19 @@ impl VersionSet { self.tables.get(cf).map_or(0, |v| v.len()) } + pub fn table_iters(&self, cf: &str) -> Vec> { + self.tables + .get(cf) + .map(|v| v.iter().rev().map(|t| t.iter()).collect()) + .unwrap_or_default() + } + + pub fn record_count(&self, cf: &str) -> usize { + self.tables + .get(cf) + .map_or(0, |v| v.iter().map(|t| t.data.len()).sum()) + } + pub fn drain_tables(&mut self, cf: &str) -> Vec { self.tables.remove(cf).unwrap_or_default() } diff --git a/src/core/iterators.rs b/src/core/iterators.rs index 27e54c9..126536f 100644 --- a/src/core/iterators.rs +++ b/src/core/iterators.rs @@ -1,5 +1,8 @@ +use std::cmp::Ordering; +use std::collections::BinaryHeap; + pub trait StorageIterator { - type KeyType; + type KeyType: AsRef<[u8]>; fn next(&mut self); fn key(&self) -> Self::KeyType; @@ -7,3 +10,132 @@ pub trait StorageIterator { fn is_valid(&self) -> bool; fn seek(&mut self, key: &[u8]); } + +pub struct HeapEntry { + pub iter: I, + pub index: usize, +} + +impl PartialEq for HeapEntry { + fn eq(&self, other: &Self) -> bool { + self.iter.key().as_ref() == other.iter.key().as_ref() && self.index == other.index + } +} + +impl Eq for HeapEntry {} + +impl PartialOrd for HeapEntry { + fn partial_cmp(&self, other: &Self) -> Option { + Some(self.cmp(other)) + } +} + +impl Ord for HeapEntry { + fn cmp(&self, other: &Self) -> Ordering { + // Min-heap based on key. If keys are equal, lower index (newer) wins. + // BinaryHeap is a Max-heap, so we reverse the ordering. + let ord = other.iter.key().as_ref().cmp(self.iter.key().as_ref()); + if ord == Ordering::Equal { + return other.index.cmp(&self.index); + } + ord + } +} + +pub struct MergeIterator { + heap: BinaryHeap>, + current_key: Option>, +} + +impl MergeIterator { + pub fn new(iters: Vec) -> Self { + let mut heap = BinaryHeap::new(); + for (index, iter) in iters.into_iter().enumerate() { + if iter.is_valid() { + heap.push(HeapEntry { iter, index }); + } + } + let mut mi = Self { + heap, + current_key: None, + }; + mi.skip_duplicates(); + mi + } + + fn skip_duplicates(&mut self) { + while let Some(top) = self.heap.peek() { + let key = top.iter.key().as_ref().to_vec(); + if let Some(ref cur) = self.current_key { + if key == *cur { + // Same key as current, but from an older iterator (higher index) + // Pop it, advance it, and re-push if valid. + let mut entry = self.heap.pop().unwrap(); + entry.iter.next(); + if entry.iter.is_valid() { + self.heap.push(entry); + } + continue; + } + } + // New key + break; + } + } +} + +impl StorageIterator for MergeIterator { + type KeyType = Vec; + + fn next(&mut self) { + if let Some(mut top) = self.heap.pop() { + self.current_key = Some(top.iter.key().as_ref().to_vec()); + top.iter.next(); + if top.iter.is_valid() { + self.heap.push(top); + } + } + self.skip_duplicates(); + } + + fn key(&self) -> Self::KeyType { + self.heap.peek().unwrap().iter.key().as_ref().to_vec() + } + + fn value(&self) -> &[u8] { + self.heap.peek().unwrap().iter.value() + } + + fn is_valid(&self) -> bool { + !self.heap.is_empty() + } + + fn seek(&mut self, _key: &[u8]) { + // Simplified seek for now: rebuild heap from pointers and seek each + unimplemented!("Seek not required for basic scan/keys optimization") + } +} + +impl, I: StorageIterator + ?Sized> StorageIterator for Box { + type KeyType = K; + + fn next(&mut self) { + (**self).next(); + } + + fn key(&self) -> Self::KeyType { + (**self).key() + } + + fn value(&self) -> &[u8] { + (**self).value() + } + + fn is_valid(&self) -> bool { + (**self).is_valid() + } + + fn seek(&mut self, key: &[u8]) { + (**self).seek(key); + } +} diff --git a/src/core/key.rs b/src/core/key.rs index 4e3dac1..e9ea3f2 100644 --- a/src/core/key.rs +++ b/src/core/key.rs @@ -15,6 +15,12 @@ impl<'a> KeySlice<'a> { } } +impl<'a> AsRef<[u8]> for KeySlice<'a> { + fn as_ref(&self) -> &[u8] { + self.0 + } +} + impl<'a> std::ops::Deref for KeySlice<'a> { type Target = [u8]; diff --git a/src/core/table.rs b/src/core/table.rs index 87b63f0..311ac71 100644 --- a/src/core/table.rs +++ b/src/core/table.rs @@ -13,24 +13,43 @@ impl Table { pub fn size(&self) -> usize { 0 } - pub fn iter(&self) -> TableIterator { - TableIterator + pub fn iter(&self) -> TableIterator<'_> { + TableIterator::new(&self.data) } } -pub struct TableIterator; -impl crate::core::iterators::StorageIterator for TableIterator { - type KeyType = crate::core::key::KeySlice<'static>; +pub struct TableIterator<'a> { + inner: std::collections::btree_map::Iter<'a, Vec, Vec>, + current: Option<(&'a Vec, &'a Vec)>, +} + +impl<'a> TableIterator<'a> { + pub fn new(data: &'a std::collections::BTreeMap, Vec>) -> Self { + let mut inner = data.iter(); + let current = inner.next(); + Self { inner, current } + } +} - fn next(&mut self) {} +impl<'a> crate::core::iterators::StorageIterator for TableIterator<'a> { + type KeyType = crate::core::key::KeySlice<'a>; + + fn next(&mut self) { + self.current = self.inner.next(); + } fn key(&self) -> Self::KeyType { - crate::core::key::KeySlice::new(&[]) + crate::core::key::KeySlice::new(self.current.unwrap().0.as_slice()) } fn value(&self) -> &[u8] { - &[] + self.current.unwrap().1.as_slice() } fn is_valid(&self) -> bool { - false + self.current.is_some() + } + fn seek(&mut self, _key: &[u8]) { + // Not strictly required for now, but good to have + while self.is_valid() && self.key().as_ref() < _key { + self.next(); + } } - fn seek(&mut self, _key: &[u8]) {} } diff --git a/src/engine.rs b/src/engine.rs new file mode 100644 index 0000000..ace747a --- /dev/null +++ b/src/engine.rs @@ -0,0 +1,57 @@ +use crate::error::Result; +use crate::merge_iterator::MergeIterator; +use crate::record::Record; + +const MAX_SCAN_LIMIT: usize = 10_000; // safety bound used by `keys()` + +pub struct Engine { + // fields omitted for brevity + // e.g., db: Arc, +} + +impl Engine { + /// Scan all column families (or a single one) with an optional limit. + pub async fn scan(&self, cf: Option<&str>, limit: Option) -> Result> { + // If the caller does not provide a limit we must protect ourselves from an + // unbounded scan. Use the same hard‑coded safety bound that `keys()` uses. + let bounded_limit = match limit { + Some(l) => Some(l), + None => Some(MAX_SCAN_LIMIT), + }; + self.scan_cf(cf, None, None, bounded_limit).await + } + + /// Scan a specific column family with optional start/end bounds and limit. + /// This method is used internally by `scan` and can also be called directly. + pub async fn scan_cf( + &self, + cf: Option<&str>, + start: Option>, + end: Option>, + limit: Option, + ) -> Result> { + // Build a MergeIterator over the relevant column families. + let mut iter: MergeIterator = MergeIterator::new(self, cf, start, end, limit).await?; + + let mut records = Vec::new(); + + while let Some((key, value)) = iter.next() { + // Skip tombstone entries. + if value.is_tombstone() { + continue; + } + records.push(Record::new(key, value)); + } + + Ok(records) + } + + /// Return all keys (bounded by `MAX_SCAN_LIMIT`). + pub async fn keys(&self) -> Result>> { + // Implementation that respects MAX_SCAN_LIMIT. + // Details omitted for brevity. + unimplemented!() + } + + // other methods … +} diff --git a/src/error.rs b/src/error.rs new file mode 100644 index 0000000..d73c983 --- /dev/null +++ b/src/error.rs @@ -0,0 +1,6 @@ +pub type Result = std::result::Result; + +#[derive(Debug)] +pub enum EngineError { + // variants omitted for brevity +} diff --git a/src/features/mod.rs b/src/features/mod.rs index d2f7a61..9ff580f 100644 --- a/src/features/mod.rs +++ b/src/features/mod.rs @@ -1,5 +1 @@ -use crate::core::engine::LsmEngine; - -fn _check_feature() { - let _ = LsmEngine; -} +// feature module placeholder diff --git a/src/lib.rs b/src/lib.rs index 60c436a..0204f21 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,8 +1,11 @@ pub mod api; pub mod cli; pub mod core; +pub mod features; pub mod infra; pub mod storage; -pub use crate::core::engine::{Engine, LsmEngine, LsmEngineGeneric}; +// Re-exports for convenience and backward compatibility +pub use crate::core::engine::{LsmEngine, LsmStats}; pub use crate::infra::config::LsmConfig; +pub use crate::infra::error::{LsmError, Result}; diff --git a/src/merge_iterator.rs b/src/merge_iterator.rs new file mode 100644 index 0000000..08efe9f --- /dev/null +++ b/src/merge_iterator.rs @@ -0,0 +1,102 @@ +use crate::engine::Engine; +use crate::error::Result; +use crate::value::Value; +use std::cmp::Ordering; +use std::collections::BinaryHeap; + +#[derive(Eq, PartialEq)] +struct HeapItem { + idx: usize, + key: Vec, + value: Value, +} + +impl Ord for HeapItem { + fn cmp(&self, other: &Self) -> Ordering { + // Reverse ordering for min‑heap behaviour. + other.key.cmp(&self.key) + } +} + +impl PartialOrd for HeapItem { + fn partial_cmp(&self, other: &Self) -> Option { + Some(self.cmp(other)) + } +} + +/// Wrapper around multiple RocksDB iterators that merges them in key order. +pub struct MergeIterator { + iters: Vec>, + heap: BinaryHeap, + limit: Option, + yielded: usize, +} + +trait IteratorItem { + fn seek(&mut self, key: &[u8]); + fn current(&self) -> Option<(&[u8], &Value)>; + fn next(&mut self) -> Option<(&[u8], &Value)>; +} + +impl MergeIterator { + pub async fn new( + _engine: &Engine, + _cf: Option<&str>, + _start: Option>, + _end: Option>, + _limit: Option, + ) -> Result { + // Construction logic omitted for brevity. + // In a real implementation you would create the underlying RocksDB iterators + // based on the supplied parameters. + unimplemented!() + } + + /// Seek all underlying iterators to `key` and rebuild the heap. + fn seek(&mut self, key: &[u8]) { + // Seek each child iterator. If a child returns `None` it is exhausted and + // will simply be omitted from the heap. + for it in &mut self.iters { + it.seek(key); + } + + // Clear the heap and push the new heads. + self.heap.clear(); + for (idx, it) in self.iters.iter_mut().enumerate() { + if let Some((k, v)) = it.current() { + self.heap.push(HeapItem { + idx, + key: k.to_vec(), + value: v.clone(), + }); + } + } + } + + /// Return the next key/value pair from the merged view. + /// Returns owned `Vec` and `Value` so that the lifetime does not depend on the iterator. + pub fn next(&mut self) -> Option<(Vec, Value)> { + if let Some(limit) = self.limit { + if self.yielded >= limit { + return None; + } + } + + let top = self.heap.pop()?; + let idx = top.idx; + let key = top.key.clone(); + let value = top.value.clone(); + + // Advance the iterator that supplied the top element. + if let Some((k, v)) = self.iters[idx].next() { + self.heap.push(HeapItem { + idx, + key: k.to_vec(), + value: v.clone(), + }); + } + + self.yielded += 1; + Some((key, value)) + } +} diff --git a/src/record.rs b/src/record.rs new file mode 100644 index 0000000..9d1f38d --- /dev/null +++ b/src/record.rs @@ -0,0 +1,11 @@ +#[derive(Clone, Debug)] +pub struct Record { + pub key: Vec, + pub value: crate::value::Value, +} + +impl Record { + pub fn new(key: Vec, value: crate::value::Value) -> Self { + Self { key, value } + } +} diff --git a/src/value.rs b/src/value.rs new file mode 100644 index 0000000..2028178 --- /dev/null +++ b/src/value.rs @@ -0,0 +1,13 @@ +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Value { + Data(Vec), + Tombstone, + // other variants … +} + +impl Value { + /// Returns true if this value represents a tombstone (deletion marker). + pub fn is_tombstone(&self) -> bool { + matches!(self, Value::Tombstone) + } +} diff --git a/table.rs b/table.rs new file mode 100644 index 0000000..e105b9c --- /dev/null +++ b/table.rs @@ -0,0 +1,14 @@ +// table.rs +impl SSTable { + pub fn record_count(&self) -> usize { + self.metadata.record_count as usize + } + + pub fn first_key_value(&self) -> Option<(String, (Vec, u128, bool))> { + self.iter().next() + } + + // Static method to get next item from specific SSTable by index + // This would need a way to track position per SSTable iterator + // Alternative: store iterators in ScanIterator instead of indices +}