//! One actor per open repository. //! //! `gigit_git::Repository` is blocking and thread-affine, so each one gets a //! thread of its own that owns it outright. Nothing else ever holds a reference //! to it, which is why there is no lock anywhere in this file. Callers get a //! [`RepositoryHandle`], which is just a channel sender and is cheap to clone. use std::path::{Path, PathBuf}; use gigit_git::{RefEntry, Repository, RepositorySummary}; use tauri::ipc::Channel; use tokio::sync::{mpsc, oneshot}; use crate::error::{Error, Result}; use crate::graph_stream::{GraphEvent, GraphStream}; /// How many requests can be in flight before callers start waiting. Requests /// are answered in microseconds to milliseconds, so this only ever fills up if /// something is badly wrong. const INBOX_CAPACITY: usize = 32; /// A request for the thread that owns a repository. /// /// Every variant carries the channel its answer goes back on, so the actor /// never needs to know who asked. enum Message { Summary { respond: oneshot::Sender>, }, References { respond: oneshot::Sender>>, }, /// Begin streaming the commit graph, replacing any stream already running. StreamGraph { channel: Channel, budget: usize, respond: oneshot::Sender>, }, /// Ask for more rows on the stream that is already running. RequestMoreRows { rows: usize, respond: oneshot::Sender>, }, } /// A cheap, clonable way to talk to one repository's thread. #[derive(Clone, Debug)] pub struct RepositoryHandle { outbox: mpsc::Sender, } impl RepositoryHandle { /// Open a repository on a thread of its own. /// /// Opening is itself blocking, so it happens on the new thread rather than /// on the caller's: the repository is owned by its thread from the very /// first moment, and the async runtime is never blocked. pub async fn open(path: PathBuf) -> Result<(Self, RepositorySummary)> { let (outbox, inbox) = mpsc::channel(INBOX_CAPACITY); let (ready, opened) = oneshot::channel(); std::thread::Builder::new() .name(format!("gigit-repository {}", path.display())) .spawn(move || run(&path, inbox, ready)) .map_err(Error::SpawnThread)?; // The thread drops `ready` without sending only if it panics. let summary = opened.await.map_err(|_| Error::WorkerStopped)??; Ok((Self { outbox }, summary)) } pub async fn summary(&self) -> Result { self.request(|respond| Message::Summary { respond }).await } pub async fn references(&self) -> Result> { self.request(|respond| Message::References { respond }) .await } /// Start streaming the commit graph down `channel`. /// /// Only one stream runs per repository; starting another replaces it, and /// the old channel simply stops receiving. pub async fn stream_graph(&self, channel: Channel, budget: usize) -> Result<()> { self.request(|respond| Message::StreamGraph { channel, budget, respond, }) .await } /// Let the running stream send another `rows` rows. pub async fn request_more_rows(&self, rows: usize) -> Result<()> { self.request(|respond| Message::RequestMoreRows { rows, respond }) .await } /// Send a request and wait for its answer. /// /// Both channels can only fail by the actor being gone, which for a handle /// held by a caller means the thread panicked or the repository was closed /// underneath them. 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)? .map_err(Error::Git) } } /// The body of a repository thread: open, report, then serve until the last /// handle is dropped and the channel closes. fn run( path: &Path, mut inbox: mpsc::Receiver, ready: oneshot::Sender>, ) { let repository = match open_and_summarise(path) { Ok(opened) => opened, Err(error) => { // Nobody waiting for the answer is not an error worth reporting — // it just means the caller gave up first. let _ = ready.send(Err(error)); return; } }; let (repository, summary) = repository; if ready.send(Ok(summary)).is_err() { // The caller went away before we finished opening, so there is nobody // left to serve. return; } // At most one graph stream per repository. Dropping it drops the walk, // which is all there is to cancelling one. let mut stream: Option = None; // `blocking_recv` is the point of running on a dedicated thread: the // blocking git work below never touches the async runtime's workers. while let Some(message) = inbox.blocking_recv() { match message { Message::Summary { respond } => { let _ = respond.send(repository.summary()); } Message::References { respond } => { let _ = respond.send(repository.references()); } Message::StreamGraph { channel, budget, respond, } => { // Replaces whatever was streaming before; the previous walk is // dropped here. stream = None; match repository.graph() { Ok(walk) => { let mut started = GraphStream::new(walk, channel, budget); if started.pump() { stream = Some(started); } let _ = respond.send(Ok(())); } Err(error) => { let _ = respond.send(Err(error)); } } } Message::RequestMoreRows { rows, respond } => { // Asking for more when nothing is streaming is not an error — // the stream may have finished between the frontend deciding // to ask and the message arriving. if let Some(running) = &mut stream { running.extend_budget(rows); if !running.pump() { stream = None; } } let _ = respond.send(Ok(())); } } } } fn open_and_summarise(path: &Path) -> Result<(Repository, RepositorySummary)> { let repository = Repository::discover(path)?; let summary = repository.summary()?; Ok((repository, summary)) }