Add eframe dependency and basic GUI structure
This commit introduces the `eframe` dependency and sets up the foundational structure for a graphical user interface. It includes the necessary imports for `egui` and prepares the application to render UI elements.
This commit is contained in:
Generated
+3290
-10
File diff suppressed because it is too large
Load Diff
@@ -10,3 +10,4 @@ tokio = { version = "1.49.0", features = ["full"] }
|
||||
zip = "8.0.0"
|
||||
sysinfo = "0.33.0"
|
||||
chrono = "0.4"
|
||||
eframe = "0.30.0"
|
||||
|
||||
+153
-141
@@ -1,24 +1,25 @@
|
||||
mod records;
|
||||
|
||||
use std::io::Read;
|
||||
use std::path::{Path, PathBuf};
|
||||
use tokio::fs;
|
||||
use zip::ZipArchive;
|
||||
use bytemuck::cast_slice;
|
||||
use crate::records::{M1Record, TickRecord, OhlcItem, DataPoint};
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::{Arc, RwLock};
|
||||
use std::time::Instant;
|
||||
use eframe::egui;
|
||||
use tokio::runtime::Runtime;
|
||||
use std::collections::HashMap;
|
||||
use sysinfo::System;
|
||||
use chrono::Local;
|
||||
use zip::ZipArchive;
|
||||
use bytemuck::cast_slice;
|
||||
use std::io::Read;
|
||||
use crate::records::{M1Record, TickRecord, OhlcItem, DataPoint};
|
||||
|
||||
// --- Data Structures ---
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct DataFile {
|
||||
pub path: PathBuf,
|
||||
pub symbol: String,
|
||||
pub year: i32,
|
||||
pub month: i32,
|
||||
pub month: i,
|
||||
}
|
||||
|
||||
struct CacheEntry<T> {
|
||||
@@ -34,10 +35,6 @@ pub struct DataServer {
|
||||
sys: RwLock<System>,
|
||||
}
|
||||
|
||||
fn log(msg: &str) {
|
||||
println!("[{}] {}", Local::now().format("%H:%M:%S%.3f"), msg);
|
||||
}
|
||||
|
||||
impl DataServer {
|
||||
pub fn new(path: impl AsRef<Path>) -> Self {
|
||||
let mut sys = System::new_all();
|
||||
@@ -52,7 +49,7 @@ impl DataServer {
|
||||
}
|
||||
|
||||
pub async fn update_symbols(&self) -> tokio::io::Result<()> {
|
||||
let mut entries = fs::read_dir(&self.base_path).await?;
|
||||
let mut entries = tokio::fs::read_dir(&self.base_path).await?;
|
||||
let mut temp_map: HashMap<String, Vec<DataFile>> = HashMap::new();
|
||||
|
||||
while let Some(entry) = entries.next_entry().await? {
|
||||
@@ -64,7 +61,6 @@ impl DataServer {
|
||||
}
|
||||
}
|
||||
|
||||
// Sort files within each symbol by year/month
|
||||
for files in temp_map.values_mut() {
|
||||
files.sort_by(|a, b| a.year.cmp(&b.year).then(a.month.cmp(&b.month)));
|
||||
}
|
||||
@@ -84,53 +80,12 @@ impl DataServer {
|
||||
fn parse_file_name(&self, path: &Path, file_name: &str) -> Option<DataFile> {
|
||||
let stem = Path::new(file_name).file_stem()?.to_str()?;
|
||||
let ext = Path::new(file_name).extension()?.to_str()?;
|
||||
|
||||
let parts: Vec<&str> = stem.split('_').collect();
|
||||
if parts.len() < 3 { return None; }
|
||||
|
||||
let month = parts.last()?.parse::<i32>().ok()?;
|
||||
let year = parts[parts.len() - 2].parse::<i32>().ok()?;
|
||||
let symbol = format!("{}.{}", parts[..parts.len() - 2].join("_"), ext);
|
||||
|
||||
Some(DataFile {
|
||||
path: path.to_path_buf(),
|
||||
symbol,
|
||||
year,
|
||||
month,
|
||||
})
|
||||
}
|
||||
|
||||
fn check_memory_and_prune(&self) {
|
||||
let mut sys = self.sys.write().unwrap();
|
||||
sys.refresh_memory();
|
||||
|
||||
let available = sys.available_memory();
|
||||
let total = sys.total_memory();
|
||||
let threshold = (total / 10).min(2 * 1024 * 1024 * 1024);
|
||||
|
||||
if available < threshold {
|
||||
self.prune_oldest();
|
||||
}
|
||||
}
|
||||
|
||||
fn prune_oldest(&self) {
|
||||
let mut m1 = self.m1_cache.write().unwrap();
|
||||
let mut tick = self.tick_cache.write().unwrap();
|
||||
|
||||
let total_before = m1.len() + tick.len();
|
||||
let mut entries: Vec<(PathBuf, Instant, bool)> = m1.iter()
|
||||
.map(|(p, e)| (p.clone(), e.last_accessed, true))
|
||||
.chain(tick.iter().map(|(p, e)| (p.clone(), e.last_accessed, false)))
|
||||
.collect();
|
||||
|
||||
entries.sort_by_key(|e| e.1);
|
||||
let to_remove = (entries.len() / 10).max(1);
|
||||
|
||||
for i in 0..to_remove.min(entries.len()) {
|
||||
let (path, _, is_m1) = &entries[i];
|
||||
if *is_m1 { m1.remove(path); } else { tick.remove(path); }
|
||||
}
|
||||
println!("\n[Cache] Memory low! Pruned {} items. ({} -> {})", to_remove, total_before, m1.len() + tick.len());
|
||||
Some(DataFile { path: path.to_path_buf(), symbol, year, month })
|
||||
}
|
||||
|
||||
pub async fn load_m1_data(&self, data_file: &DataFile) -> anyhow::Result<Arc<Vec<DataPoint<OhlcItem>>>> {
|
||||
@@ -141,61 +96,21 @@ impl DataServer {
|
||||
return Ok(Arc::clone(&entry.data));
|
||||
}
|
||||
}
|
||||
|
||||
self.check_memory_and_prune();
|
||||
|
||||
let raw_data = self.read_from_zip::<M1Record>(&data_file.path).await?;
|
||||
let mapped_data: Vec<DataPoint<OhlcItem>> = raw_data.into_iter()
|
||||
.map(|r| r.to_ohlc())
|
||||
.collect();
|
||||
|
||||
let mapped_data: Vec<DataPoint<OhlcItem>> = raw_data.into_iter().map(|r| r.to_ohlc()).collect();
|
||||
let shared_data = Arc::new(mapped_data);
|
||||
|
||||
let mut cache = self.m1_cache.write().unwrap();
|
||||
cache.insert(data_file.path.clone(), CacheEntry {
|
||||
data: Arc::clone(&shared_data),
|
||||
last_accessed: Instant::now(),
|
||||
});
|
||||
|
||||
Ok(shared_data)
|
||||
}
|
||||
|
||||
pub async fn load_tick_data(&self, data_file: &DataFile) -> anyhow::Result<Arc<Vec<TickRecord>>> {
|
||||
{
|
||||
let mut cache = self.tick_cache.write().unwrap();
|
||||
if let Some(entry) = cache.get_mut(&data_file.path) {
|
||||
entry.last_accessed = Instant::now();
|
||||
return Ok(Arc::clone(&entry.data));
|
||||
}
|
||||
}
|
||||
|
||||
self.check_memory_and_prune();
|
||||
let data = self.read_from_zip::<TickRecord>(&data_file.path).await?;
|
||||
let shared_data = Arc::new(data);
|
||||
|
||||
let mut cache = self.tick_cache.write().unwrap();
|
||||
cache.insert(data_file.path.clone(), CacheEntry {
|
||||
data: Arc::clone(&shared_data),
|
||||
last_accessed: Instant::now(),
|
||||
});
|
||||
|
||||
cache.insert(data_file.path.clone(), CacheEntry { data: Arc::clone(&shared_data), last_accessed: Instant::now() });
|
||||
Ok(shared_data)
|
||||
}
|
||||
|
||||
async fn read_from_zip<R: bytemuck::Pod>(&self, path: &Path) -> anyhow::Result<Vec<R>> {
|
||||
let file = std::fs::File::open(path)?;
|
||||
let mut archive = ZipArchive::new(file)?;
|
||||
|
||||
for i in 0..archive.len() {
|
||||
let mut file = archive.by_index(i)?;
|
||||
if file.name().ends_with(".bin") {
|
||||
let size = file.size() as usize;
|
||||
let record_size = std::mem::size_of::<R>();
|
||||
if size % record_size != 0 {
|
||||
return Err(anyhow::anyhow!("Invalid file size for record type"));
|
||||
}
|
||||
|
||||
let mut bin_content = vec![0u8; size];
|
||||
let mut bin_content = vec![0u8; file.size() as usize];
|
||||
file.read_exact(&mut bin_content)?;
|
||||
let records: &[R] = cast_slice(&bin_content);
|
||||
return Ok(records.to_vec());
|
||||
@@ -203,56 +118,153 @@ impl DataServer {
|
||||
}
|
||||
Err(anyhow::anyhow!("No .bin file found"))
|
||||
}
|
||||
|
||||
pub async fn process_symbol_m1(&self, symbol: &str) -> anyhow::Result<()> {
|
||||
let files = {
|
||||
let symbols = self.symbols.read().unwrap();
|
||||
symbols.get(symbol).cloned().ok_or_else(|| anyhow::anyhow!("Symbol not found"))?
|
||||
};
|
||||
|
||||
log(&format!(">>> Loading Symbol: {} ({} files)", symbol, files.len()));
|
||||
|
||||
for file in files {
|
||||
if file.symbol.ends_with(".m1") {
|
||||
let _ = self.load_m1_data(&file).await?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> anyhow::Result<()> {
|
||||
let server = Arc::new(DataServer::new(r"\\COFFEE\TickData\Pepperstone"));
|
||||
// --- EGui App ---
|
||||
|
||||
log("Indexing symbols...");
|
||||
server.update_symbols().await?;
|
||||
struct DataViewerApp {
|
||||
server: Arc<DataServer>,
|
||||
rt: Runtime,
|
||||
symbols: Vec<String>,
|
||||
selected_symbol: Option<String>,
|
||||
loaded_data: Option<Arc<Vec<DataPoint<OhlcItem>>>>,
|
||||
tx: std::sync::mpsc::Sender<Arc<Vec<DataPoint<OhlcItem>>>>,
|
||||
rx: std::sync::mpsc::Receiver<Arc<Vec<DataPoint<OhlcItem>>>>,
|
||||
is_loading: bool,
|
||||
status_msg: String,
|
||||
}
|
||||
|
||||
let symbols = server.enumerate_symbols();
|
||||
log(&format!("Found {} symbols.", symbols.len()));
|
||||
impl DataViewerApp {
|
||||
fn new(cc: &eframe::CreationContext<'_>) -> Self {
|
||||
let server = Arc::new(DataServer::new(r"\\COFFEE\TickData\Pepperstone"));
|
||||
let rt = Runtime::new().unwrap();
|
||||
|
||||
let start_all = Instant::now();
|
||||
let mut set = tokio::task::JoinSet::new();
|
||||
let max_concurrency = 8; // Lower concurrency since we load full symbols now
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
for symbol in symbols {
|
||||
let s = Arc::clone(&server);
|
||||
// Initial symbol load
|
||||
let server_clone = Arc::clone(&server);
|
||||
let symbols = rt.block_on(async {
|
||||
server_clone.update_symbols().await.ok();
|
||||
server_clone.enumerate_symbols()
|
||||
});
|
||||
|
||||
while set.len() >= max_concurrency {
|
||||
set.join_next().await;
|
||||
Self {
|
||||
server,
|
||||
rt,
|
||||
symbols,
|
||||
selected_symbol: None,
|
||||
loaded_data: None,
|
||||
tx,
|
||||
rx,
|
||||
is_loading: false,
|
||||
status_msg: "Ready".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
set.spawn(async move {
|
||||
let _ = s.process_symbol_m1(&symbol).await;
|
||||
fn load_symbol_async(&mut self, symbol: String) {
|
||||
self.is_loading = true;
|
||||
self.status_msg = format!("Loading {}...", symbol);
|
||||
self.loaded_data = None;
|
||||
|
||||
let server = Arc::clone(&self.server);
|
||||
let tx = self.tx.clone();
|
||||
let ctx = self.rt.handle().clone();
|
||||
|
||||
self.rt.spawn(async move {
|
||||
let files = {
|
||||
let symbols = server.symbols.read().unwrap();
|
||||
symbols.get(&symbol).cloned().unwrap_or_default()
|
||||
};
|
||||
|
||||
let mut all_data = Vec::new();
|
||||
for file in files {
|
||||
if file.symbol.ends_with(".m1") {
|
||||
if let Ok(data) = server.load_m1_data(&file).await {
|
||||
all_data.extend((*data).clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
let _ = tx.send(Arc::new(all_data));
|
||||
});
|
||||
}
|
||||
|
||||
while let Some(_) = set.join_next().await {}
|
||||
|
||||
log(&format!("Finished loading all symbols in {:?}", start_all.elapsed()));
|
||||
|
||||
let m1_count = server.m1_cache.read().unwrap().len();
|
||||
log(&format!("Final Cache State: {} M1 files in RAM.", m1_count));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
impl eframe::App for DataViewerApp {
|
||||
fn update(&mut self, ctx: &egui::Context, _frame: &mut eframe::Frame) {
|
||||
if self.is_loading {
|
||||
ctx.request_repaint();
|
||||
}
|
||||
|
||||
// Poll for finished loading tasks
|
||||
if let Ok(data) = self.rx.try_recv() {
|
||||
self.loaded_data = Some(data);
|
||||
self.is_loading = false;
|
||||
self.status_msg = format!("Loaded {} records", self.loaded_data.as_ref().unwrap().len());
|
||||
}
|
||||
|
||||
let mut symbol_to_load = None;
|
||||
|
||||
egui::SidePanel::left("symbol_panel").show(ctx, |ui| {
|
||||
ui.heading("Symbols");
|
||||
egui::ScrollArea::vertical().show(ui, |ui| {
|
||||
for symbol in &self.symbols {
|
||||
let is_selected = self.selected_symbol.as_ref() == Some(symbol);
|
||||
if ui.selectable_label(is_selected, symbol).clicked() && !is_selected {
|
||||
symbol_to_load = Some(symbol.clone());
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
if let Some(symbol) = symbol_to_load {
|
||||
self.selected_symbol = Some(symbol.clone());
|
||||
self.load_symbol_async(symbol);
|
||||
}
|
||||
|
||||
egui::CentralPanel::default().show(ctx, |ui| {
|
||||
if let Some(symbol) = &self.selected_symbol {
|
||||
ui.heading(format!("Data: {}", symbol));
|
||||
if self.is_loading {
|
||||
ui.add(egui::Spinner::new());
|
||||
ui.label("Streaming from network share...");
|
||||
} else if let Some(data) = &self.loaded_data {
|
||||
ui.label(format!("Successfully loaded {} records.", data.len()));
|
||||
|
||||
// Show a small sample
|
||||
egui::ScrollArea::vertical().show(ui, |ui| {
|
||||
for i in 0..data.len().min(10000) {
|
||||
let d = &data[i];
|
||||
// Copy fields to local variables to avoid unaligned references
|
||||
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));
|
||||
}
|
||||
});
|
||||
}
|
||||
} else {
|
||||
ui.label("Bitte ein Symbol auswählen.");
|
||||
}
|
||||
});
|
||||
|
||||
egui::TopBottomPanel::bottom("status_bar").show(ctx, |ui| {
|
||||
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));
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
fn main() -> eframe::Result {
|
||||
let native_options = eframe::NativeOptions::default();
|
||||
eframe::run_native(
|
||||
"Rust Data Viewer",
|
||||
native_options,
|
||||
Box::new(|cc| Ok(Box::new(DataViewerApp::new(cc)))),
|
||||
)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user