use super::{DiffPage, PreparedSnapshot};
use crate::{
cancellation::{AgentCancellation, AgentCancellationHandle},
diff_review::{self, CommentChange},
};
use std::sync::Arc;
use std::{
path::{Path, PathBuf},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
struct Outcome {
snapshot: Result<Arc<PreparedSnapshot>, String>,
mutation: Option<Result<(), String>>,
}
#[derive(Default)]
pub(crate) struct DiffWorker {
task: Option<JoinHandle<Outcome>>,
cancel: Option<AgentCancellationHandle>,
refresh_invalidated: bool,
cwd: PathBuf,
next_refresh: Option<Instant>,
cache: Arc<std::sync::Mutex<diff_review::SnapshotCache>>,
projections: Arc<std::sync::Mutex<super::projection::ProjectionCache>>,
}
impl DiffWorker {
pub(crate) fn is_pending(&self) -> bool {
self.task.is_some()
}
pub(crate) fn poll(&mut self, page: &mut DiffPage, cwd: &Path, active: bool) -> bool {
let mut changed = false;
if !active
&& !page.saving
&& let Some(cancel) = &self.cancel
{
cancel.cancel();
self.refresh_invalidated = true;
}
if self.cwd != cwd {
if let Some(cancel) = &self.cancel {
cancel.cancel();
self.refresh_invalidated = true;
}
self.next_refresh = None;
}
if self.task.as_ref().is_some_and(|task| task.is_finished()) {
let result = self.task.take().expect("finished worker").join();
self.cancel = None;
let refresh_invalidated = std::mem::take(&mut self.refresh_invalidated);
if self.cwd == cwd {
match result {
Ok(outcome) => {
if outcome.mutation.is_none()
&& !refresh_invalidated
&& outcome.snapshot.is_ok()
&& page.dialog.is_none()
{
page.error = None;
}
if let Some(result) = outcome.mutation {
page.saving = false;
match result {
Ok(()) => {
page.dialog = None;
page.error = None;
}
Err(error) => page.error = Some(error),
}
}
if !refresh_invalidated {
match outcome.snapshot {
Ok(snapshot) => page.apply_snapshot(snapshot),
Err(error) if active => page.error = Some(error),
Err(_) => {}
}
}
}
Err(_) => {
page.saving = false;
page.error = Some("Diff worker stopped unexpectedly".into());
}
}
changed = true;
}
self.next_refresh = (!refresh_invalidated).then(|| {
Instant::now()
+ if page
.snapshot
.as_ref()
.is_some_and(|s| s.highlighting_pending)
{
Duration::from_millis(16)
} else {
Duration::from_secs(1)
}
});
}
if self.task.is_some() {
return changed;
}
if self.cwd != cwd {
self.cwd = cwd.to_path_buf();
*page = DiffPage::default();
self.next_refresh = None;
changed = true;
}
if !active && page.pending.is_none() {
return changed;
}
if page.pending.is_none() && self.next_refresh.is_some_and(|next| Instant::now() < next) {
return changed;
}
let change: Option<CommentChange> = page.pending.take();
let mutation_requested = change.is_some();
let path = cwd.to_path_buf();
let (token, cancel) = AgentCancellation::default().child_token();
let previous = page.snapshot.clone();
let selected_path = previous
.as_ref()
.and_then(|snapshot| snapshot.files.get(page.selected))
.map(|file| file.path.clone());
let cache = Arc::clone(&self.cache);
let projections = Arc::clone(&self.projections);
match thread::Builder::new()
.name("magi-diff-review".into())
.spawn(move || {
let mutation = change.map(|change| {
if token.is_canceled() {
Err("Comment update canceled".into())
} else {
diff_review::save_comment(&path, change, &token)
}
});
let snapshot = if token.is_canceled() {
Err("Diff refresh canceled".into())
} else {
cache
.lock()
.map_err(|_| "Diff cache lock poisoned".to_owned())
.and_then(|mut cache| {
let source = diff_review::load_snapshot(&path, &token, &mut cache)?;
token.check().map_err(|error| error.to_string())?;
let mut projections = projections
.lock()
.map_err(|_| "Diff projection cache lock poisoned".to_owned())?;
let snapshot = PreparedSnapshot::prepare(
source,
previous.as_deref(),
selected_path.as_deref(),
&token,
&mut projections,
)
.map_err(|error| error.to_string())?;
token.check().map_err(|error| error.to_string())?;
Ok(Arc::new(snapshot))
})
};
Outcome { snapshot, mutation }
}) {
Ok(task) => {
self.task = Some(task);
self.cancel = Some(cancel);
}
Err(error) => {
page.saving = false;
page.error = Some(format!("Cannot start diff worker: {error}"));
self.next_refresh = Some(Instant::now() + Duration::from_secs(1));
changed = true;
}
}
changed || mutation_requested
}
}
impl Drop for DiffWorker {
fn drop(&mut self) {
if let Some(cancel) = &self.cancel {
cancel.cancel();
}
if self.task.as_ref().is_some_and(|task| task.is_finished()) {
let _ = self.task.take().expect("finished worker").join();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn partial_snapshot_schedules_automatic_continuation() {
let source = diff_review::ReviewSnapshot {
root: "/review".into(),
files: vec![diff_review::ReviewFile {
path: "file.rs".into(),
original: "fn original() {}".into(),
current: "fn current() {}".into(),
rows: vec![],
notice: None,
}],
comments: vec![],
};
let token = AgentCancellation::default();
let mut worker = DiffWorker::default();
worker.cwd = "/review".into();
let snapshot = PreparedSnapshot::with_budget(
source,
None,
None,
&crate::rendering::highlight::SourceHighlightBudget {
deadline: Instant::now(),
cancellation: &token,
},
&mut worker.projections.lock().unwrap(),
1,
)
.unwrap();
worker.task = Some(thread::spawn(move || Outcome {
snapshot: Ok(Arc::new(snapshot)),
mutation: None,
}));
wait_for_task(&worker);
let mut page = DiffPage::default();
assert!(worker.poll(&mut page, Path::new("/review"), true));
assert!(page.snapshot.as_ref().unwrap().highlighting_pending);
assert!(
worker
.next_refresh
.unwrap()
.saturating_duration_since(Instant::now())
<= Duration::from_millis(16)
);
assert!(worker.task.is_none());
}
fn isolate_review_storage(
env: &crate::test_support::env::EnvGuard,
) -> (tempfile::TempDir, crate::test_support::env::SavedEnvVar<'_>) {
let directory = tempfile::tempdir().unwrap();
let saved_home = env.save("MC_HOME");
env.set_var("MC_HOME", directory.path().canonicalize().unwrap());
(directory, saved_home)
}
fn wait_for_task(worker: &DiffWorker) {
let deadline = Instant::now() + Duration::from_secs(10);
while !worker.task.as_ref().unwrap().is_finished() {
assert!(Instant::now() < deadline, "diff worker did not finish");
thread::yield_now();
}
}
fn repository() -> tempfile::TempDir {
let directory = tempfile::tempdir().unwrap();
assert!(
std::process::Command::new("git")
.args(["init", "--quiet"])
.current_dir(directory.path())
.status()
.unwrap()
.success()
);
std::fs::write(directory.path().join("first.txt"), "first change\n").unwrap();
directory
}
fn blocked_refresh(worker: &mut DiffWorker, cwd: &Path) -> std::sync::mpsc::SyncSender<()> {
let (release, wait) = std::sync::mpsc::sync_channel(1);
let (token, cancel) = AgentCancellation::default().child_token();
worker.cwd = cwd.to_path_buf();
worker.cancel = Some(cancel);
worker.task = Some(thread::spawn(move || {
wait.recv_timeout(Duration::from_secs(10)).unwrap();
Outcome {
snapshot: Err(token.check().unwrap_err().to_string()),
mutation: None,
}
}));
release
}
#[test]
fn first_open_loads_without_canceling_its_token() {
let env = crate::test_support::env::env_lock();
let (_storage, _saved_home) = isolate_review_storage(&env);
let directory = repository();
let mut worker = DiffWorker::default();
let mut page = DiffPage::default();
worker.poll(&mut page, directory.path(), false);
assert!(!worker.is_pending());
worker.poll(&mut page, directory.path(), true);
wait_for_task(&worker);
worker.poll(&mut page, directory.path(), true);
assert!(page.error.is_none(), "{:?}", page.error);
assert_eq!(page.snapshot.unwrap().files[0].path, "first.txt");
}
#[test]
fn reopening_retries_canceled_first_refresh_without_showing_an_error() {
let env = crate::test_support::env::env_lock();
let (_storage, _saved_home) = isolate_review_storage(&env);
let directory = repository();
let mut worker = DiffWorker::default();
let mut page = DiffPage::default();
let release = blocked_refresh(&mut worker, directory.path());
worker.poll(&mut page, directory.path(), false);
worker.poll(&mut page, directory.path(), true);
release.send(()).unwrap();
wait_for_task(&worker);
worker.poll(&mut page, directory.path(), true);
assert!(page.error.is_none(), "{:?}", page.error);
assert!(worker.is_pending(), "reopening must retry immediately");
wait_for_task(&worker);
worker.poll(&mut page, directory.path(), true);
assert!(page.error.is_none(), "{:?}", page.error);
assert_eq!(page.snapshot.unwrap().files[0].path, "first.txt");
}
#[test]
fn current_refresh_errors_are_not_hidden() {
let cwd = PathBuf::from("/review");
for error in ["Git review timed out", "prompt canceled"] {
let mut worker = DiffWorker::default();
worker.cwd = cwd.clone();
worker.task = Some(thread::spawn(move || Outcome {
snapshot: Err(error.into()),
mutation: None,
}));
let mut page = DiffPage::default();
wait_for_task(&worker);
worker.poll(&mut page, &cwd, true);
assert_eq!(page.error.as_deref(), Some(error));
assert!(!worker.is_pending());
}
}
#[test]
fn cwd_change_during_first_refresh_loads_the_new_repository() {
let env = crate::test_support::env::env_lock();
let (_storage, _saved_home) = isolate_review_storage(&env);
let first = repository();
let second = repository();
let mut worker = DiffWorker::default();
let mut page = DiffPage::default();
let release = blocked_refresh(&mut worker, first.path());
worker.poll(&mut page, second.path(), true);
release.send(()).unwrap();
wait_for_task(&worker);
worker.poll(&mut page, second.path(), true);
assert!(page.error.is_none());
assert!(worker.is_pending());
wait_for_task(&worker);
worker.poll(&mut page, second.path(), true);
assert!(page.error.is_none(), "{:?}", page.error);
assert_eq!(
page.snapshot.unwrap().root.canonicalize().unwrap(),
second.path().canonicalize().unwrap()
);
}
#[test]
fn refreshed_pane_focus_changes_borders_not_empty_backgrounds() {
use crate::tui::state::{MissionControlState, RightColumnTab, TuiFocusPane};
use ratatui::{Terminal, backend::TestBackend};
let env = crate::test_support::env::env_lock();
let (_storage, _saved_home) = isolate_review_storage(&env);
let directory = repository();
let mut worker = DiffWorker::default();
let mut state = MissionControlState::default();
state.select_right_column_tab(RightColumnTab::Diff);
worker.poll(&mut state.diff, directory.path(), true);
wait_for_task(&worker);
worker.poll(&mut state.diff, directory.path(), true);
assert!(state.diff.snapshot.is_some());
let mut terminal = Terminal::new(TestBackend::new(100, 12)).unwrap();
for focus in 1..=3 {
state.focus_pane = TuiFocusPane::Transcript;
state.diff.focus = focus;
terminal
.draw(|frame| super::super::render::draw(frame, frame.area(), &state))
.unwrap();
let focused = terminal.backend().buffer().clone();
state.focus_pane = TuiFocusPane::Prompt;
terminal
.draw(|frame| super::super::render::draw(frame, frame.area(), &state))
.unwrap();
let unfocused = terminal.backend().buffer();
for (index, column) in super::super::render::columns(focused.area)
.iter()
.enumerate()
{
for x in column.x + 1..column.right() - 1 {
assert_eq!(
focused[(x, 10)],
unfocused[(x, 10)],
"pane {index}, focus {focus}"
);
}
let border = (column.x, 5);
if index + 1 == focus {
assert_ne!(focused[border].fg, unfocused[border].fg);
} else {
assert_eq!(focused[border], unfocused[border]);
}
}
}
}
#[test]
fn retained_mutation_ack_is_applied_even_after_tab_is_closed() {
let cwd = PathBuf::from("/review");
let mut worker = DiffWorker::default();
worker.cwd = cwd.clone();
worker.task = Some(thread::spawn(|| Outcome {
snapshot: Err("refresh failed after save".into()),
mutation: Some(Ok(())),
}));
wait_for_task(&worker);
let mut page = DiffPage {
saving: true,
..Default::default()
};
assert!(worker.poll(&mut page, &cwd, false));
assert!(!page.saving);
assert!(!worker.is_pending());
}
#[test]
fn failed_comment_write_keeps_draft_for_retry() {
let cwd = PathBuf::from("/review");
let mut worker = DiffWorker::default();
worker.cwd = cwd.clone();
worker.task = Some(thread::spawn(|| Outcome {
snapshot: Ok(Arc::new(PreparedSnapshot::new(
crate::diff_review::ReviewSnapshot {
root: "/review".into(),
files: vec![],
comments: vec![],
},
None,
))),
mutation: Some(Err("disk is read-only".into())),
}));
wait_for_task(&worker);
let mut page = DiffPage {
saving: true,
..Default::default()
};
page.dialog = Some(super::super::CommentDialog::from_comment(
&crate::diff_review::ReviewComment {
id: "comment".into(),
path: "file.rs".into(),
side: crate::diff_review::ReviewSide::Changed,
line: 1,
anchor: "source".into(),
text: "keep this draft".into(),
stale: false,
resolved: false,
},
));
assert!(worker.poll(&mut page, &cwd, true));
assert!(!page.saving);
assert_eq!(page.error.as_deref(), Some("disk is read-only"));
assert_eq!(
page.dialog.as_ref().unwrap().editor.text(),
"keep this draft"
);
assert!(!worker.is_pending());
}
}