//! The actor that owns every open repository. //! //! It owns the map, so the map needs no lock either. Tauri commands hold a //! [`WorkspaceHandle`] and nothing else. use std::collections::HashMap; use std::future::Future; use std::path::PathBuf; use gigit_git::{RefEntry, RepositorySummary}; use serde::{Deserialize, Serialize}; use tauri::ipc::Channel; use tokio::sync::{mpsc, oneshot}; use crate::error::{Error, Result}; use crate::graph_stream::GraphEvent; use crate::recents::{RecentRepository, RecentsHandle}; use crate::repository_actor::RepositoryHandle; const INBOX_CAPACITY: usize = 32; /// Identifies an open repository. /// /// It is the canonical path to the `.git` directory, which makes opening the /// same repository twice — by a different spelling of the path, or from a /// subdirectory — resolve to the one that is already open. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct RepositoryId(String); impl RepositoryId { fn for_summary(summary: &RepositorySummary) -> Self { let path = summary .git_dir .canonicalize() .unwrap_or_else(|_| summary.git_dir.clone()); Self(path.to_string_lossy().into_owned()) } } impl std::fmt::Display for RepositoryId { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter.write_str(&self.0) } } /// What the frontend gets back when a repository opens. #[derive(Clone, Debug, Serialize)] #[serde(rename_all = "camelCase")] pub struct OpenedRepository { pub id: RepositoryId, pub summary: RepositorySummary, } enum Message { Open { path: PathBuf, respond: oneshot::Sender>, }, Close { id: RepositoryId, respond: oneshot::Sender, }, Lookup { id: RepositoryId, respond: oneshot::Sender>, }, } #[derive(Clone, Debug)] pub struct WorkspaceHandle { outbox: mpsc::Sender, } impl WorkspaceHandle { /// Build the handle and the actor's future. /// /// The future is returned rather than spawned so the caller decides which /// runtime runs it — Tauri's in the app, the test runtime in tests. /// /// The recents handle is taken here rather than reached for in the command /// layer, because remembering a repository is part of opening one and /// commands are meant to be free of decisions like that. pub fn new(recents: RecentsHandle) -> (Self, impl Future + Send + 'static) { let (outbox, inbox) = mpsc::channel(INBOX_CAPACITY); (Self { outbox }, run(inbox, recents)) } /// Open a repository, or return the one already open for that path. pub async fn open(&self, path: PathBuf) -> Result { // Two layers of failure: reaching the actor, and the open itself. self.request(|respond| Message::Open { path, respond }) .await? } /// Close a repository. Returns whether it was open to begin with. pub async fn close(&self, id: RepositoryId) -> Result { self.request(|respond| Message::Close { id, respond }).await } pub async fn references(&self, id: RepositoryId) -> Result> { self.repository(id).await?.references().await } pub async fn summary(&self, id: RepositoryId) -> Result { self.repository(id).await?.summary().await } pub async fn stream_graph( &self, id: RepositoryId, channel: Channel, budget: usize, ) -> Result<()> { self.repository(id) .await? .stream_graph(channel, budget) .await } pub async fn request_more_rows(&self, id: RepositoryId, rows: usize) -> Result<()> { self.repository(id).await?.request_more_rows(rows).await } async fn repository(&self, id: RepositoryId) -> Result { let found = self .request(|respond| Message::Lookup { id: id.clone(), respond, }) .await?; found.ok_or_else(|| Error::UnknownRepository(id.to_string())) } async fn request(&self, build: impl FnOnce(oneshot::Sender) -> Message) -> Result { let (respond, answer) = oneshot::channel(); self.outbox .send(build(respond)) .await .map_err(|_| Error::WorkerStopped)?; answer.await.map_err(|_| Error::WorkerStopped) } } async fn run(mut inbox: mpsc::Receiver, recents: RecentsHandle) { let mut open: HashMap = HashMap::new(); while let Some(message) = inbox.recv().await { match message { // Opening is awaited here rather than spawned, so two opens are // served one after the other. Opening is fast enough that the // simplicity is worth more than the concurrency. Message::Open { path, respond } => { let opened = open_repository(&mut open, path).await; // Only a repository that actually opened is worth remembering, // and failing to remember one is never worth failing the open. if let Ok(opened) = &opened { let entry = RecentRepository::new(opened.id.clone(), &opened.summary); let _ = recents.record(entry).await; } let _ = respond.send(opened); } Message::Close { id, respond } => { // Dropping the last handle closes the actor's channel, which // ends its loop and drops the repository with it. let _ = respond.send(open.remove(&id).is_some()); } Message::Lookup { id, respond } => { let _ = respond.send(open.get(&id).cloned()); } } } } async fn open_repository( open: &mut HashMap, path: PathBuf, ) -> Result { let (handle, summary) = RepositoryHandle::open(path).await?; let id = RepositoryId::for_summary(&summary); // The id is only knowable after opening, so a repository that is already // open costs one short-lived thread: dropping the surplus handle closes // its channel, which ends that thread's loop straight away. Avoiding it // would mean resolving the git directory separately from opening it, // which is not worth the complexity for a path taken once per repository. let handle = open.entry(id.clone()).or_insert(handle); let summary = handle.summary().await?; Ok(OpenedRepository { id, summary }) }