From 147b4b745c7e370b5a6d205810b6c9d58f4e2c63 Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Mon, 16 Feb 2026 15:27:26 +0100 Subject: [PATCH] Refactor DataServer to use generic cache --- src/main.rs | 244 +++++++++++++++++++++++++++++++++---------------- src/records.rs | 15 ++- 2 files changed, 175 insertions(+), 84 deletions(-) diff --git a/src/main.rs b/src/main.rs index f0f62e3..00fb272 100644 --- a/src/main.rs +++ b/src/main.rs @@ -27,23 +27,21 @@ struct CacheEntry { last_accessed: Instant, } -pub struct DataServer { +pub struct DataServer { base_path: PathBuf, symbols: RwLock>>, - m1_cache: RwLock>>>, - tick_cache: RwLock>>, + cache: RwLock>>, sys: RwLock, } -impl DataServer { +impl DataServer { pub fn new(path: impl AsRef) -> Self { let mut sys = System::new_all(); sys.refresh_memory(); Self { base_path: path.as_ref().to_path_buf(), symbols: RwLock::new(HashMap::new()), - m1_cache: RwLock::new(HashMap::new()), - tick_cache: RwLock::new(HashMap::new()), + cache: RwLock::new(HashMap::new()), sys: RwLock::new(sys), } } @@ -55,11 +53,13 @@ impl DataServer { while let Some(entry) = entries.next_entry().await? { let path = entry.path(); if let Some(file_name) = path.file_name().and_then(|n| n.to_str()) { - if let Some(data_file) = self.parse_file_name(&path, file_name) { - temp_map - .entry(data_file.symbol.clone()) - .or_default() - .push(data_file); + if file_name.ends_with(".m1") || file_name.ends_with(".tick") { + if let Some(data_file) = self.parse_file_name(&path, file_name) { + temp_map + .entry(data_file.symbol.clone()) + .or_default() + .push(data_file); + } } } } @@ -98,22 +98,26 @@ impl DataServer { }) } - pub async fn load_m1_data( - &self, - data_file: &DataFile, - ) -> anyhow::Result>>> { + pub async fn load_data(&self, data_file: &DataFile) -> anyhow::Result>> + where + R: crate::records::SourceRecord, + crate::records::DataPoint: Into, + { { - let mut cache = self.m1_cache.write().unwrap(); + let mut cache = self.cache.write().unwrap(); if let Some(entry) = cache.get_mut(&data_file.path) { entry.last_accessed = Instant::now(); return Ok(Arc::clone(&entry.data)); } } - let raw_data = self.read_from_zip::(&data_file.path).await?; - let mapped_data: Vec> = - raw_data.into_iter().map(|r| r.to_ohlc()).collect(); + let raw_data = self.read_from_zip::(&data_file.path).await?; + let mapped_data: Vec = raw_data + .into_iter() + .map(|r| r.to_data_point().into()) + .collect(); + let shared_data = Arc::new(mapped_data); - let mut cache = self.m1_cache.write().unwrap(); + let mut cache = self.cache.write().unwrap(); cache.insert( data_file.path.clone(), CacheEntry { @@ -146,38 +150,59 @@ impl DataServer { // --- EGui App --- +// --- App Types --- + +enum LoadedData { + None, + M1(Vec>), + Ticks(Vec>), +} + +enum DataChunk { + M1(Arc>>), + Ticks(Arc>>), +} + struct DataViewerApp { - server: Arc, + m1_server: Arc>>, + tick_server: Arc>>, rt: Runtime, symbols: Vec, selected_symbol: Option, - loaded_data: Vec>, - tx: std::sync::mpsc::Sender>>>>, - rx: std::sync::mpsc::Receiver>>>>, + loaded_data: LoadedData, + tx: std::sync::mpsc::Sender>, + rx: std::sync::mpsc::Receiver>, is_loading: bool, status_msg: String, } impl DataViewerApp { - fn new(cc: &eframe::CreationContext<'_>) -> Self { - let server = Arc::new(DataServer::new(r"\\COFFEE\TickData\Pepperstone")); + fn new(_cc: &eframe::CreationContext<'_>) -> Self { + let base_path = r"\\COFFEE\TickData\Pepperstone"; + let m1_server = Arc::new(DataServer::>::new(base_path)); + let tick_server = Arc::new(DataServer::>::new(base_path)); let rt = Runtime::new().unwrap(); let (tx, rx) = std::sync::mpsc::channel(); - // Initial symbol load - let server_clone = Arc::clone(&server); + // Initial symbol load for both servers + let m1_clone = Arc::clone(&m1_server); + let tick_clone = Arc::clone(&tick_server); let symbols = rt.block_on(async { - server_clone.update_symbols().await.ok(); - server_clone.enumerate_symbols() + let _ = tokio::join!( + m1_clone.update_symbols(), + tick_clone.update_symbols() + ); + m1_clone.enumerate_symbols() }); Self { - server, + m1_server, + tick_server, rt, symbols, selected_symbol: None, - loaded_data: Vec::new(), + loaded_data: LoadedData::None, tx, rx, is_loading: false, @@ -188,44 +213,68 @@ impl DataViewerApp { fn load_symbol_async(&mut self, symbol: String) { self.is_loading = true; self.status_msg = format!("Loading {}...", symbol); - self.loaded_data.clear(); + self.loaded_data = if symbol.ends_with(".m1") { + LoadedData::M1(Vec::new()) + } else { + LoadedData::Ticks(Vec::new()) + }; - let server = Arc::clone(&self.server); + let m1_server = Arc::clone(&self.m1_server); + let tick_server = Arc::clone(&self.tick_server); let tx = self.tx.clone(); + let sym_clone = symbol.clone(); self.rt.spawn(async move { - let files: Vec = { - let symbols = server.symbols.read().unwrap(); - symbols.get(&symbol).cloned().unwrap_or_default() - }; - - let mut next_load: Option>>>>> = None; - - for i in 0..files.len() { - let file = &files[i]; - if !file.symbol.ends_with(".m1") { - continue; - } - - let data_res: anyhow::Result>>> = if let Some(handle) = next_load.take() { - handle.await.map_err(|e| anyhow::anyhow!("Join error: {}", e)).and_then(|res| res) - } else { - server.load_m1_data(file).await + if sym_clone.ends_with(".m1") { + let files = { + let symbols = m1_server.symbols.read().unwrap(); + symbols.get(&sym_clone).cloned().unwrap_or_default() }; + let mut next_load: Option>>>>> = None; - if i + 1 < files.len() { - let next_file = files[i + 1].clone(); - let server_clone = Arc::clone(&server); - next_load = Some(tokio::spawn(async move { - server_clone.load_m1_data(&next_file).await - })); + for i in 0..files.len() { + let data_res = if let Some(handle) = next_load.take() { + handle.await.map_err(|e| anyhow::anyhow!("Join error: {}", e)).and_then(|res| res) + } else { + m1_server.load_data::(&files[i]).await + }; + + if i + 1 < files.len() { + let next_file = files[i + 1].clone(); + let s = Arc::clone(&m1_server); + next_load = Some(tokio::spawn(async move { s.load_data::(&next_file).await })); + } + + if let Ok(data) = data_res { + let _ = tx.send(Some(DataChunk::M1(data))); + } } + } else { + let files = { + let symbols = tick_server.symbols.read().unwrap(); + symbols.get(&sym_clone).cloned().unwrap_or_default() + }; + let mut next_load: Option>>>>> = None; - if let Ok(data) = data_res { - let _ = tx.send(Some(data)); + for i in 0..files.len() { + let data_res = if let Some(handle) = next_load.take() { + handle.await.map_err(|e| anyhow::anyhow!("Join error: {}", e)).and_then(|res| res) + } else { + tick_server.load_data::(&files[i]).await + }; + + if i + 1 < files.len() { + let next_file = files[i + 1].clone(); + let s = Arc::clone(&tick_server); + next_load = Some(tokio::spawn(async move { s.load_data::(&next_file).await })); + } + + if let Ok(data) = data_res { + let _ = tx.send(Some(DataChunk::Ticks(data))); + } } } - let _ = tx.send(None); // Finished + let _ = tx.send(None); }); } } @@ -240,7 +289,18 @@ impl eframe::App for DataViewerApp { let mut received_any = false; while let Ok(msg) = self.rx.try_recv() { if let Some(chunk) = msg { - self.loaded_data.extend_from_slice(&chunk); + match chunk { + DataChunk::M1(data) => { + if let LoadedData::M1(vec) = &mut self.loaded_data { + vec.extend_from_slice(&data); + } + } + DataChunk::Ticks(data) => { + if let LoadedData::Ticks(vec) = &mut self.loaded_data { + vec.extend_from_slice(&data); + } + } + } received_any = true; } else { self.is_loading = false; @@ -248,7 +308,12 @@ impl eframe::App for DataViewerApp { } if received_any { - self.status_msg = format!("Loaded {} records", self.loaded_data.len()); + let count = match &self.loaded_data { + LoadedData::M1(v) => v.len(), + LoadedData::Ticks(v) => v.len(), + _ => 0, + }; + self.status_msg = format!("Loaded {} records", count); } let mut symbol_to_load = None; @@ -280,23 +345,41 @@ impl eframe::App for DataViewerApp { }); } - ui.label(format!("Successfully loaded {} records.", self.loaded_data.len())); - - // Show a small sample - egui::ScrollArea::vertical().show(ui, |ui| { - let data = &self.loaded_data; - for i in 0..data.len().min(10000) { - let d = &data[i]; - let open = d.data.open; - let high = d.data.high; - let low = d.data.low; - let close = d.data.close; - ui.label(format!( - "{}: O:{:.5} H:{:.5} L:{:.5} C:{:.5}", - i, open, high, low, close - )); + match &self.loaded_data { + LoadedData::M1(data) => { + ui.label(format!("Successfully loaded {} OHLC records.", data.len())); + egui::ScrollArea::vertical().show(ui, |ui| { + for i in 0..data.len().min(10000) { + let d = &data[i]; + let open = d.data.open; + let high = d.data.high; + let low = d.data.low; + let close = d.data.close; + ui.label(format!( + "{}: O:{:.5} H:{:.5} L:{:.5} C:{:.5}", + i, open, high, low, close + )); + } + }); } - }); + LoadedData::Ticks(data) => { + ui.label(format!("Successfully loaded {} Tick records.", data.len())); + egui::ScrollArea::vertical().show(ui, |ui| { + for i in 0..data.len().min(10000) { + let d = &data[i]; + let ask = d.data.ask; + let bid = d.data.bid; + ui.label(format!( + "{}: A:{:.5} B:{:.5}", + i, ask, bid + )); + } + }); + } + LoadedData::None => { + ui.label("Warte auf Daten..."); + } + } } else { ui.label("Bitte ein Symbol auswählen."); } @@ -306,8 +389,9 @@ impl eframe::App for DataViewerApp { ui.horizontal(|ui| { ui.label(&self.status_msg); ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| { - let m1_cached = self.server.m1_cache.read().unwrap().len(); - ui.label(format!("Cache: {} files", m1_cached)); + let m1_cached = self.m1_server.cache.read().unwrap().len(); + let tick_cached = self.tick_server.cache.read().unwrap().len(); + ui.label(format!("Cache: M1:{}, Ticks:{}", m1_cached, tick_cached)); }); }); }); diff --git a/src/records.rs b/src/records.rs index a26e59e..c234063 100644 --- a/src/records.rs +++ b/src/records.rs @@ -41,8 +41,14 @@ pub struct DataPoint { pub data: T, } -impl M1Record { - pub fn to_ohlc(self) -> DataPoint { +pub trait SourceRecord: Pod + Send + 'static { + type Mapped: Pod + Send + 'static; + fn to_data_point(self) -> DataPoint; +} + +impl SourceRecord for M1Record { + type Mapped = OhlcItem; + fn to_data_point(self) -> DataPoint { DataPoint { time: self.time, data: OhlcItem { @@ -56,8 +62,9 @@ impl M1Record { } } -impl TickRecord { - pub fn to_data_point(self) -> DataPoint { +impl SourceRecord for TickRecord { + type Mapped = TickRecord; + fn to_data_point(self) -> DataPoint { DataPoint { time: self.time, data: self,