//! Streaming graph rows to the frontend. //! //! Rows are pushed rather than returned, and the push is bounded by a budget //! the frontend tops up as it scrolls. That is credit-based flow control: the //! first rows leave for the webview the moment they exist, without waiting to //! be asked, and a million-commit repository still never floods the IPC //! boundary because nothing is sent beyond what has been asked for. //! //! A plain request/response would have been simpler, but it costs a round trip //! before anything can be drawn and gives the backend no way to speak on its //! own — which is what pushing new commits into an open view will need later. use gigit_git::{GraphRow, GraphWalk}; use serde::Serialize; use tauri::ipc::Channel; /// Rows per event. Large enough that the per-message overhead is irrelevant, /// small enough that a big budget arrives as several paintable batches rather /// than one long freeze. const ROWS_PER_EVENT: usize = 200; /// What the frontend receives on the channel. #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase", tag = "event", content = "data")] pub enum GraphEvent { Rows { /// Row index of the first row in this batch, so the frontend can place /// them without tracking how many it has seen. start: u32, rows: Vec, }, /// The history is fully drawn; no more events will arrive. Exhausted { total: u32, /// Whether the walk had to fall back on approximate ordering. degraded: bool, }, Failed { message: String, }, } /// An open graph stream: a walk, where to send it, and how much has been asked /// for but not yet sent. pub struct GraphStream { walk: GraphWalk, channel: Channel, budget: usize, sent: u32, } impl GraphStream { pub fn new(walk: GraphWalk, channel: Channel, budget: usize) -> Self { Self { walk, channel, budget, sent: 0, } } /// Raise the budget, as the frontend scrolls towards the end of what it has. pub fn extend_budget(&mut self, rows: usize) { self.budget = self.budget.saturating_add(rows); } /// Send whatever the budget allows. /// /// Returns whether the stream is still alive. A finished history, a failed /// walk, or a channel the webview has dropped all end it — and returning /// `false` is what makes the caller drop the walk, which is the whole of /// cancellation. #[must_use] pub fn pump(&mut self) -> bool { while self.budget > 0 { let wanted = self.budget.min(ROWS_PER_EVENT); let rows = match self.walk.next_chunk(wanted) { Ok(rows) => rows, Err(error) => { // A failure to send here means nobody is listening anyway. let _ = self.channel.send(GraphEvent::Failed { message: error.to_string(), }); return false; } }; if rows.is_empty() { break; } self.budget -= rows.len(); let start = self.sent; self.sent += u32::try_from(rows.len()).unwrap_or(u32::MAX); if self.channel.send(GraphEvent::Rows { start, rows }).is_err() { // The webview closed the channel: stop walking rather than // computing rows nobody will ever see. return false; } } if self.walk.is_exhausted() { let _ = self.channel.send(GraphEvent::Exhausted { total: self.sent, degraded: self.walk.is_degraded(), }); return false; } true } }