diff options
| author | gennyble <gen@nyble.dev> | 2026-09-24 01:52:48 -0500 |
|---|---|---|
| committer | gennyble <gen@nyble.dev> | 2026-09-24 01:52:48 -0500 |
| commit | f2a9bc39cb9aa8a9d9e186bfd71dfd2405fe1ab6 (patch) | |
| tree | 9a5fdc35930ce7d0a039fe5eae154555cf0231fb /splorm/src/connector | |
| parent | 32cbf8332bdaafae8c0934ff2ab45048bc8f2361 (diff) | |
splorm, meow
Diffstat (limited to 'splorm/src/connector')
| -rw-r--r-- | splorm/src/connector/gatherer.rs | 177 | ||||
| -rw-r--r-- | splorm/src/connector/mod.rs | 5 |
2 files changed, 182 insertions, 0 deletions
diff --git a/splorm/src/connector/gatherer.rs b/splorm/src/connector/gatherer.rs new file mode 100644 index 0000000..bb69196 --- /dev/null +++ b/splorm/src/connector/gatherer.rs @@ -0,0 +1,177 @@ +use std::{ops::DerefMut, sync::Arc, thread::JoinHandle}; + +use camino::{Utf8Path, Utf8PathBuf}; +use gatherer_client::{Client, Connection, Error as GathererError, GeneratedGraphs, GraphBounds}; +use timeparse::{Duration, OffsetDateTime}; +use tokio::sync::{ + Mutex, + mpsc::{Receiver, Sender}, +}; + +struct CacheItem<T: Clone> { + item: Mutex<Option<T>>, + time: Mutex<OffsetDateTime>, + timeout: Duration, +} + +impl<T: Clone> CacheItem<T> { + pub fn new(timeout: Duration) -> Self { + Self { + item: Mutex::new(None), + time: Mutex::new(timeparse::now_utc()), + timeout, + } + } + + pub fn new_with_item(item: T, timeout: Duration) -> Self { + Self { + item: Mutex::new(Some(item)), + time: Mutex::new(timeparse::now_utc()), + timeout, + } + } + + pub async fn get(&self) -> Option<T> { + if self.duration_since_cache().await > self.timeout { + let _ = self.item.lock().await.take(); + None + } else { + let item = self.item.lock().await.clone(); + item + } + } + + pub async fn store(&self, item: T) { + let mut item_lock = self.item.lock().await; + let mut time_lock = self.time.lock().await; + + *item_lock.deref_mut() = Some(item); + *time_lock.deref_mut() = timeparse::now_utc(); + } + + pub async fn duration_since_cache(&self) -> Duration { + let time_lock = self.time.lock().await; + *time_lock - timeparse::now_utc() + } +} + +pub struct GathererThread { + cmd_tx: Sender<GathererCommand>, + data_rx: Receiver<GathererData>, + handle: JoinHandle<()>, + + regular_graph_cache: Arc<CacheItem<GeneratedGraphs>>, +} + +impl GathererThread { + pub fn spawn(socket_path: &Utf8Path) -> Result<Self, GathererError> { + let (cmd_tx, cmd_rx) = tokio::sync::mpsc::channel(8); + let (data_tx, data_rx) = tokio::sync::mpsc::channel(8); + let client = Client::new(socket_path.to_owned()); + let connection = client.connect()?; + + let handle = std::thread::spawn(|| handler(connection, cmd_rx, data_tx)); + + Ok(Self { + cmd_tx, + data_rx, + handle, + + regular_graph_cache: Arc::new(CacheItem::new(Duration::minutes(5))), + }) + } + + pub fn data_rx(&mut self) -> &mut Receiver<GathererData> { + &mut self.data_rx + } + + pub async fn get_messsage(&mut self) -> GathererData { + match self.data_rx.recv().await { + None => panic!("gatherer connector channel closed"), + Some(data) => data, + } + } + + pub fn cmd_tx(&mut self) -> Sender<GathererCommand> { + self.cmd_tx.clone() + } + + pub async fn regular_graphs(&mut self) -> GeneratedGraphs { + if let Some(graphs) = self.regular_graph_cache.get().await { + graphs + } else { + self.cmd_tx + .send(GathererCommand::RegularGraphs) + .await + .unwrap(); + + match self.get_messsage().await { + GathererData::Graphs(graphs) => { + self.regular_graph_cache.store(graphs.clone()).await; + graphs + } + _ => panic!(), + } + } + } +} + +pub enum GathererCommand { + GetBounds, + GenerateGraphs(Utf8PathBuf), + RegularGraphs, +} + +pub enum GathererData { + Bounds(GraphBounds), + Graphs(GeneratedGraphs), + Error(gatherer_client::Error), +} + +fn handler( + mut conn: Connection, + mut cmd_rx: Receiver<GathererCommand>, + data_tx: Sender<GathererData>, +) { + macro_rules! unwrap_or_senderr { + ($operation:expr) => { + match $operation { + Err(err) => { + data_tx.blocking_send(GathererData::Error(err)).unwrap(); + continue; + } + Ok(ok) => ok, + } + }; + } + + loop { + while let Some(cmd) = cmd_rx.blocking_recv() { + match cmd { + GathererCommand::GetBounds => { + let bounds = unwrap_or_senderr!(conn.graph_bounds()); + data_tx.blocking_send(GathererData::Bounds(bounds)).unwrap(); + } + GathererCommand::GenerateGraphs(path) => { + let graphs = unwrap_or_senderr!(conn.generate_graphs(path)); + data_tx.blocking_send(GathererData::Graphs(graphs)).unwrap(); + } + GathererCommand::RegularGraphs => { + let graphs = unwrap_or_senderr!(conn.regular_graphs()); + data_tx.blocking_send(GathererData::Graphs(graphs)).unwrap(); + } + } + } + } +} + +/* +Currently there is an encoding issue in blobby, it seems, that is causing the +message from gatherer to be corrupted, and we're getting a panic or just incorrect +paths during decoding. + +ALSO: +we want to make this connector have a fallible conneciton. It does not need to +always be connected, and it will keep trying to connect/try to reconnect if +it disconnects. +*/ diff --git a/splorm/src/connector/mod.rs b/splorm/src/connector/mod.rs new file mode 100644 index 0000000..a2fbc85 --- /dev/null +++ b/splorm/src/connector/mod.rs @@ -0,0 +1,5 @@ +//! Connectors to external things + +mod gatherer; + +pub use gatherer::{GathererCommand, GathererData, GathererThread}; |
