about summary refs log tree commit diff
path: root/splorm/src/connector
diff options
context:
space:
mode:
authorgennyble <gen@nyble.dev>2026-09-24 01:52:48 -0500
committergennyble <gen@nyble.dev>2026-09-24 01:52:48 -0500
commitf2a9bc39cb9aa8a9d9e186bfd71dfd2405fe1ab6 (patch)
tree9a5fdc35930ce7d0a039fe5eae154555cf0231fb /splorm/src/connector
parent32cbf8332bdaafae8c0934ff2ab45048bc8f2361 (diff)
splorm, meow
Diffstat (limited to 'splorm/src/connector')
-rw-r--r--splorm/src/connector/gatherer.rs177
-rw-r--r--splorm/src/connector/mod.rs5
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};