Refactor DataServer to use generic cache
This commit is contained in:
+141
-57
@@ -27,23 +27,21 @@ struct CacheEntry<T> {
|
|||||||
last_accessed: Instant,
|
last_accessed: Instant,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct DataServer {
|
pub struct DataServer<T: bytemuck::Pod + Send + 'static> {
|
||||||
base_path: PathBuf,
|
base_path: PathBuf,
|
||||||
symbols: RwLock<HashMap<String, Vec<DataFile>>>,
|
symbols: RwLock<HashMap<String, Vec<DataFile>>>,
|
||||||
m1_cache: RwLock<HashMap<PathBuf, CacheEntry<DataPoint<OhlcItem>>>>,
|
cache: RwLock<HashMap<PathBuf, CacheEntry<T>>>,
|
||||||
tick_cache: RwLock<HashMap<PathBuf, CacheEntry<TickRecord>>>,
|
|
||||||
sys: RwLock<System>,
|
sys: RwLock<System>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DataServer {
|
impl<T: bytemuck::Pod + Send + 'static> DataServer<T> {
|
||||||
pub fn new(path: impl AsRef<Path>) -> Self {
|
pub fn new(path: impl AsRef<Path>) -> Self {
|
||||||
let mut sys = System::new_all();
|
let mut sys = System::new_all();
|
||||||
sys.refresh_memory();
|
sys.refresh_memory();
|
||||||
Self {
|
Self {
|
||||||
base_path: path.as_ref().to_path_buf(),
|
base_path: path.as_ref().to_path_buf(),
|
||||||
symbols: RwLock::new(HashMap::new()),
|
symbols: RwLock::new(HashMap::new()),
|
||||||
m1_cache: RwLock::new(HashMap::new()),
|
cache: RwLock::new(HashMap::new()),
|
||||||
tick_cache: RwLock::new(HashMap::new()),
|
|
||||||
sys: RwLock::new(sys),
|
sys: RwLock::new(sys),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -55,6 +53,7 @@ impl DataServer {
|
|||||||
while let Some(entry) = entries.next_entry().await? {
|
while let Some(entry) = entries.next_entry().await? {
|
||||||
let path = entry.path();
|
let path = entry.path();
|
||||||
if let Some(file_name) = path.file_name().and_then(|n| n.to_str()) {
|
if let Some(file_name) = path.file_name().and_then(|n| n.to_str()) {
|
||||||
|
if file_name.ends_with(".m1") || file_name.ends_with(".tick") {
|
||||||
if let Some(data_file) = self.parse_file_name(&path, file_name) {
|
if let Some(data_file) = self.parse_file_name(&path, file_name) {
|
||||||
temp_map
|
temp_map
|
||||||
.entry(data_file.symbol.clone())
|
.entry(data_file.symbol.clone())
|
||||||
@@ -63,6 +62,7 @@ impl DataServer {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
for files in temp_map.values_mut() {
|
for files in temp_map.values_mut() {
|
||||||
files.sort_by(|a, b| a.year.cmp(&b.year).then(a.month.cmp(&b.month)));
|
files.sort_by(|a, b| a.year.cmp(&b.year).then(a.month.cmp(&b.month)));
|
||||||
@@ -98,22 +98,26 @@ impl DataServer {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn load_m1_data(
|
pub async fn load_data<R>(&self, data_file: &DataFile) -> anyhow::Result<Arc<Vec<T>>>
|
||||||
&self,
|
where
|
||||||
data_file: &DataFile,
|
R: crate::records::SourceRecord,
|
||||||
) -> anyhow::Result<Arc<Vec<DataPoint<OhlcItem>>>> {
|
crate::records::DataPoint<R::Mapped>: Into<T>,
|
||||||
{
|
{
|
||||||
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) {
|
if let Some(entry) = cache.get_mut(&data_file.path) {
|
||||||
entry.last_accessed = Instant::now();
|
entry.last_accessed = Instant::now();
|
||||||
return Ok(Arc::clone(&entry.data));
|
return Ok(Arc::clone(&entry.data));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let raw_data = self.read_from_zip::<M1Record>(&data_file.path).await?;
|
let raw_data = self.read_from_zip::<R>(&data_file.path).await?;
|
||||||
let mapped_data: Vec<DataPoint<OhlcItem>> =
|
let mapped_data: Vec<T> = raw_data
|
||||||
raw_data.into_iter().map(|r| r.to_ohlc()).collect();
|
.into_iter()
|
||||||
|
.map(|r| r.to_data_point().into())
|
||||||
|
.collect();
|
||||||
|
|
||||||
let shared_data = Arc::new(mapped_data);
|
let shared_data = Arc::new(mapped_data);
|
||||||
let mut cache = self.m1_cache.write().unwrap();
|
let mut cache = self.cache.write().unwrap();
|
||||||
cache.insert(
|
cache.insert(
|
||||||
data_file.path.clone(),
|
data_file.path.clone(),
|
||||||
CacheEntry {
|
CacheEntry {
|
||||||
@@ -146,38 +150,59 @@ impl DataServer {
|
|||||||
|
|
||||||
// --- EGui App ---
|
// --- EGui App ---
|
||||||
|
|
||||||
|
// --- App Types ---
|
||||||
|
|
||||||
|
enum LoadedData {
|
||||||
|
None,
|
||||||
|
M1(Vec<DataPoint<OhlcItem>>),
|
||||||
|
Ticks(Vec<DataPoint<TickRecord>>),
|
||||||
|
}
|
||||||
|
|
||||||
|
enum DataChunk {
|
||||||
|
M1(Arc<Vec<DataPoint<OhlcItem>>>),
|
||||||
|
Ticks(Arc<Vec<DataPoint<TickRecord>>>),
|
||||||
|
}
|
||||||
|
|
||||||
struct DataViewerApp {
|
struct DataViewerApp {
|
||||||
server: Arc<DataServer>,
|
m1_server: Arc<DataServer<DataPoint<OhlcItem>>>,
|
||||||
|
tick_server: Arc<DataServer<DataPoint<TickRecord>>>,
|
||||||
rt: Runtime,
|
rt: Runtime,
|
||||||
symbols: Vec<String>,
|
symbols: Vec<String>,
|
||||||
selected_symbol: Option<String>,
|
selected_symbol: Option<String>,
|
||||||
loaded_data: Vec<DataPoint<OhlcItem>>,
|
loaded_data: LoadedData,
|
||||||
tx: std::sync::mpsc::Sender<Option<Arc<Vec<DataPoint<OhlcItem>>>>>,
|
tx: std::sync::mpsc::Sender<Option<DataChunk>>,
|
||||||
rx: std::sync::mpsc::Receiver<Option<Arc<Vec<DataPoint<OhlcItem>>>>>,
|
rx: std::sync::mpsc::Receiver<Option<DataChunk>>,
|
||||||
is_loading: bool,
|
is_loading: bool,
|
||||||
status_msg: String,
|
status_msg: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DataViewerApp {
|
impl DataViewerApp {
|
||||||
fn new(cc: &eframe::CreationContext<'_>) -> Self {
|
fn new(_cc: &eframe::CreationContext<'_>) -> Self {
|
||||||
let server = Arc::new(DataServer::new(r"\\COFFEE\TickData\Pepperstone"));
|
let base_path = r"\\COFFEE\TickData\Pepperstone";
|
||||||
|
let m1_server = Arc::new(DataServer::<DataPoint<OhlcItem>>::new(base_path));
|
||||||
|
let tick_server = Arc::new(DataServer::<DataPoint<TickRecord>>::new(base_path));
|
||||||
let rt = Runtime::new().unwrap();
|
let rt = Runtime::new().unwrap();
|
||||||
|
|
||||||
let (tx, rx) = std::sync::mpsc::channel();
|
let (tx, rx) = std::sync::mpsc::channel();
|
||||||
|
|
||||||
// Initial symbol load
|
// Initial symbol load for both servers
|
||||||
let server_clone = Arc::clone(&server);
|
let m1_clone = Arc::clone(&m1_server);
|
||||||
|
let tick_clone = Arc::clone(&tick_server);
|
||||||
let symbols = rt.block_on(async {
|
let symbols = rt.block_on(async {
|
||||||
server_clone.update_symbols().await.ok();
|
let _ = tokio::join!(
|
||||||
server_clone.enumerate_symbols()
|
m1_clone.update_symbols(),
|
||||||
|
tick_clone.update_symbols()
|
||||||
|
);
|
||||||
|
m1_clone.enumerate_symbols()
|
||||||
});
|
});
|
||||||
|
|
||||||
Self {
|
Self {
|
||||||
server,
|
m1_server,
|
||||||
|
tick_server,
|
||||||
rt,
|
rt,
|
||||||
symbols,
|
symbols,
|
||||||
selected_symbol: None,
|
selected_symbol: None,
|
||||||
loaded_data: Vec::new(),
|
loaded_data: LoadedData::None,
|
||||||
tx,
|
tx,
|
||||||
rx,
|
rx,
|
||||||
is_loading: false,
|
is_loading: false,
|
||||||
@@ -188,44 +213,68 @@ impl DataViewerApp {
|
|||||||
fn load_symbol_async(&mut self, symbol: String) {
|
fn load_symbol_async(&mut self, symbol: String) {
|
||||||
self.is_loading = true;
|
self.is_loading = true;
|
||||||
self.status_msg = format!("Loading {}...", symbol);
|
self.status_msg = format!("Loading {}...", symbol);
|
||||||
self.loaded_data.clear();
|
self.loaded_data = if symbol.ends_with(".m1") {
|
||||||
|
LoadedData::M1(Vec::new())
|
||||||
let server = Arc::clone(&self.server);
|
} else {
|
||||||
let tx = self.tx.clone();
|
LoadedData::Ticks(Vec::new())
|
||||||
|
|
||||||
self.rt.spawn(async move {
|
|
||||||
let files: Vec<DataFile> = {
|
|
||||||
let symbols = server.symbols.read().unwrap();
|
|
||||||
symbols.get(&symbol).cloned().unwrap_or_default()
|
|
||||||
};
|
};
|
||||||
|
|
||||||
|
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 {
|
||||||
|
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<tokio::task::JoinHandle<anyhow::Result<Arc<Vec<DataPoint<OhlcItem>>>>>> = None;
|
let mut next_load: Option<tokio::task::JoinHandle<anyhow::Result<Arc<Vec<DataPoint<OhlcItem>>>>>> = None;
|
||||||
|
|
||||||
for i in 0..files.len() {
|
for i in 0..files.len() {
|
||||||
let file = &files[i];
|
let data_res = if let Some(handle) = next_load.take() {
|
||||||
if !file.symbol.ends_with(".m1") {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
let data_res: anyhow::Result<Arc<Vec<DataPoint<OhlcItem>>>> = if let Some(handle) = next_load.take() {
|
|
||||||
handle.await.map_err(|e| anyhow::anyhow!("Join error: {}", e)).and_then(|res| res)
|
handle.await.map_err(|e| anyhow::anyhow!("Join error: {}", e)).and_then(|res| res)
|
||||||
} else {
|
} else {
|
||||||
server.load_m1_data(file).await
|
m1_server.load_data::<M1Record>(&files[i]).await
|
||||||
};
|
};
|
||||||
|
|
||||||
if i + 1 < files.len() {
|
if i + 1 < files.len() {
|
||||||
let next_file = files[i + 1].clone();
|
let next_file = files[i + 1].clone();
|
||||||
let server_clone = Arc::clone(&server);
|
let s = Arc::clone(&m1_server);
|
||||||
next_load = Some(tokio::spawn(async move {
|
next_load = Some(tokio::spawn(async move { s.load_data::<M1Record>(&next_file).await }));
|
||||||
server_clone.load_m1_data(&next_file).await
|
|
||||||
}));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Ok(data) = data_res {
|
if let Ok(data) = data_res {
|
||||||
let _ = tx.send(Some(data));
|
let _ = tx.send(Some(DataChunk::M1(data)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let _ = tx.send(None); // Finished
|
} else {
|
||||||
|
let files = {
|
||||||
|
let symbols = tick_server.symbols.read().unwrap();
|
||||||
|
symbols.get(&sym_clone).cloned().unwrap_or_default()
|
||||||
|
};
|
||||||
|
let mut next_load: Option<tokio::task::JoinHandle<anyhow::Result<Arc<Vec<DataPoint<TickRecord>>>>>> = None;
|
||||||
|
|
||||||
|
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::<TickRecord>(&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::<TickRecord>(&next_file).await }));
|
||||||
|
}
|
||||||
|
|
||||||
|
if let Ok(data) = data_res {
|
||||||
|
let _ = tx.send(Some(DataChunk::Ticks(data)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let _ = tx.send(None);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -240,7 +289,18 @@ impl eframe::App for DataViewerApp {
|
|||||||
let mut received_any = false;
|
let mut received_any = false;
|
||||||
while let Ok(msg) = self.rx.try_recv() {
|
while let Ok(msg) = self.rx.try_recv() {
|
||||||
if let Some(chunk) = msg {
|
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;
|
received_any = true;
|
||||||
} else {
|
} else {
|
||||||
self.is_loading = false;
|
self.is_loading = false;
|
||||||
@@ -248,7 +308,12 @@ impl eframe::App for DataViewerApp {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if received_any {
|
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;
|
let mut symbol_to_load = None;
|
||||||
@@ -280,11 +345,10 @@ impl eframe::App for DataViewerApp {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
ui.label(format!("Successfully loaded {} records.", self.loaded_data.len()));
|
match &self.loaded_data {
|
||||||
|
LoadedData::M1(data) => {
|
||||||
// Show a small sample
|
ui.label(format!("Successfully loaded {} OHLC records.", data.len()));
|
||||||
egui::ScrollArea::vertical().show(ui, |ui| {
|
egui::ScrollArea::vertical().show(ui, |ui| {
|
||||||
let data = &self.loaded_data;
|
|
||||||
for i in 0..data.len().min(10000) {
|
for i in 0..data.len().min(10000) {
|
||||||
let d = &data[i];
|
let d = &data[i];
|
||||||
let open = d.data.open;
|
let open = d.data.open;
|
||||||
@@ -297,6 +361,25 @@ impl eframe::App for DataViewerApp {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
}
|
||||||
|
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 {
|
} else {
|
||||||
ui.label("Bitte ein Symbol auswählen.");
|
ui.label("Bitte ein Symbol auswählen.");
|
||||||
}
|
}
|
||||||
@@ -306,8 +389,9 @@ impl eframe::App for DataViewerApp {
|
|||||||
ui.horizontal(|ui| {
|
ui.horizontal(|ui| {
|
||||||
ui.label(&self.status_msg);
|
ui.label(&self.status_msg);
|
||||||
ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| {
|
ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| {
|
||||||
let m1_cached = self.server.m1_cache.read().unwrap().len();
|
let m1_cached = self.m1_server.cache.read().unwrap().len();
|
||||||
ui.label(format!("Cache: {} files", m1_cached));
|
let tick_cached = self.tick_server.cache.read().unwrap().len();
|
||||||
|
ui.label(format!("Cache: M1:{}, Ticks:{}", m1_cached, tick_cached));
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
+11
-4
@@ -41,8 +41,14 @@ pub struct DataPoint<T> {
|
|||||||
pub data: T,
|
pub data: T,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl M1Record {
|
pub trait SourceRecord: Pod + Send + 'static {
|
||||||
pub fn to_ohlc(self) -> DataPoint<OhlcItem> {
|
type Mapped: Pod + Send + 'static;
|
||||||
|
fn to_data_point(self) -> DataPoint<Self::Mapped>;
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SourceRecord for M1Record {
|
||||||
|
type Mapped = OhlcItem;
|
||||||
|
fn to_data_point(self) -> DataPoint<OhlcItem> {
|
||||||
DataPoint {
|
DataPoint {
|
||||||
time: self.time,
|
time: self.time,
|
||||||
data: OhlcItem {
|
data: OhlcItem {
|
||||||
@@ -56,8 +62,9 @@ impl M1Record {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TickRecord {
|
impl SourceRecord for TickRecord {
|
||||||
pub fn to_data_point(self) -> DataPoint<TickRecord> {
|
type Mapped = TickRecord;
|
||||||
|
fn to_data_point(self) -> DataPoint<TickRecord> {
|
||||||
DataPoint {
|
DataPoint {
|
||||||
time: self.time,
|
time: self.time,
|
||||||
data: self,
|
data: self,
|
||||||
|
|||||||
Reference in New Issue
Block a user