use super::super::*;
use super::{
MissionControlApp, PendingRewind, RewindRequestKind, RewindWorker, RewindWorkerDelivery,
RewindWorkerOutcome,
};
use crate::checkpoints::RewindMode;
fn execute_conversation_rewind(
paths: &crate::config::McPaths,
cwd: &std::path::Path,
store: &crate::checkpoints::CheckpointStore,
plan: &crate::checkpoints::RestorePlan,
session: Option<&Session>,
mode: RewindMode,
) -> anyhow::Result<TuiRewindWorkerResult> {
let source = session.ok_or_else(|| anyhow::anyhow!("rewind requires a persisted session"))?;
anyhow::ensure!(source.id() == plan.session_id, "rewind session changed");
let manager = SessionManager::new(paths.sessions.clone());
let fork = source.fork_before_prompt(&manager, plan.target_turn)?;
let mut text = format!(
"Rewind fork: {}. Original session: {}.",
fork.id(),
source.id()
);
if let Err(error) = store.copy_before_turn(source.id(), fork.id(), plan.target_turn) {
anyhow::bail!(
"fork {} preserved; files unchanged; checkpoint copy failed: {}",
fork.id(),
error
);
}
if mode == RewindMode::Both {
let execution = store.execute_plan(plan, cwd);
let payload = crate::checkpoints::rewind_event_payload(plan, &execution);
let record = record_session_event(Some(&fork), cwd, SessionEventKind::Rewind, payload);
text.push('\n');
text.push_str(&crate::shell::format_rewind_execution_for_tui(&execution));
if let Err(error) = record {
anyhow::bail!(
"fork {} preserved; files may have changed; session recording failed: {}",
fork.id(),
error
);
}
}
record_session_event(
Some(&fork),
cwd,
SessionEventKind::Diagnostic,
serde_json::json!({"level": "info", "message": text}),
)
.map_err(|error| {
anyhow::anyhow!(
"fork {} preserved; result recording failed: {}",
fork.id(),
error
)
})?;
Ok(TuiRewindWorkerResult::Forked {
session_id: fork.id().to_owned(),
session: fork,
text,
})
}
#[derive(Debug)]
enum RewindJob {
Changes,
Selection,
Plan {
target_turn: u64,
dry_run: bool,
mode: RewindMode,
},
Execute {
plan: crate::checkpoints::RestorePlan,
session: Option<Session>,
mode: RewindMode,
},
}
fn bounded_rewind_diagnostic(error: impl std::fmt::Display) -> String {
const MAX_REWIND_DIAGNOSTIC_CHARS: usize = 500;
crate::output::sanitize_display_text(&error.to_string())
.chars()
.take(MAX_REWIND_DIAGNOSTIC_CHARS)
.collect()
}
fn rewind_worker_result_kind(result: &TuiRewindWorkerResult) -> RewindRequestKind {
match result {
TuiRewindWorkerResult::Changes { .. } => RewindRequestKind::Changes,
TuiRewindWorkerResult::Plan { .. } => RewindRequestKind::Plan,
TuiRewindWorkerResult::Execution { .. } => RewindRequestKind::Execute,
TuiRewindWorkerResult::Forked { .. } => RewindRequestKind::Execute,
}
}
fn run_rewind_job(
paths: crate::config::McPaths,
session_id: String,
cwd: PathBuf,
job: RewindJob,
cancel: &AtomicBool,
) -> Option<Result<TuiRewindWorkerResult, String>> {
if cancel.load(Ordering::SeqCst) {
return None;
}
let store = crate::checkpoints::CheckpointStore::from_paths(&paths);
match job {
RewindJob::Changes => {
let changes = store.changes(&session_id);
if cancel.load(Ordering::SeqCst) {
return None;
}
Some(Ok(TuiRewindWorkerResult::Changes {
text: crate::shell::format_changes_for_tui(&changes),
plans: Vec::new(),
}))
}
RewindJob::Selection => {
let read = store.read_records(&session_id);
let changes = store.changes_from_read(&read);
if cancel.load(Ordering::SeqCst) {
return None;
}
let text = crate::shell::format_changes_for_tui(&changes);
let mut plans = Vec::new();
let session =
match SessionManager::new(paths.sessions.clone()).open_existing(&session_id) {
Ok(session) => session,
Err(error) => return Some(Err(bounded_rewind_diagnostic(error))),
};
let prompts = match session.rewind_prompts() {
Ok(prompts) => prompts,
Err(error) => return Some(Err(bounded_rewind_diagnostic(error))),
};
let first = prompts.len().saturating_sub(100);
for prompt in prompts.into_iter().skip(first) {
if cancel.load(Ordering::SeqCst) {
return None;
}
let plan = store.plan_rewind_from_read(&session_id, &cwd, prompt.turn, &read);
let text = format!(
"Latest 100 prompts; use /rewind --to <turn> for an earlier retained prompt.\n{}",
crate::shell::format_rewind_prompt_text(&plan, &prompt.text)
);
plans.push(TuiRewindPlanPreview { plan, text });
}
Some(Ok(TuiRewindWorkerResult::Changes { text, plans }))
}
RewindJob::Plan {
target_turn,
dry_run,
mode,
} => {
let plan = store.plan_rewind(&session_id, &cwd, target_turn);
let text = match SessionManager::new(paths.sessions.clone())
.open_existing(&session_id)
.and_then(|session| crate::shell::format_rewind_prompt_preview(&session, &plan))
{
Ok(text) => text,
Err(error) => return Some(Err(bounded_rewind_diagnostic(error))),
};
if cancel.load(Ordering::SeqCst) {
return None;
}
Some(Ok(TuiRewindWorkerResult::Plan {
text,
plan,
dry_run,
mode,
}))
}
RewindJob::Execute {
plan,
session,
mode,
} => {
if mode != RewindMode::Files {
return Some(
execute_conversation_rewind(
&paths,
&cwd,
&store,
&plan,
session.as_ref(),
mode,
)
.map_err(bounded_rewind_diagnostic),
);
}
let execution = store.execute_plan(&plan, &cwd);
let session_record_error = record_session_event(
session.as_ref(),
&cwd,
SessionEventKind::Rewind,
crate::checkpoints::rewind_event_payload(&plan, &execution),
)
.err()
.map(bounded_rewind_diagnostic);
let mut text = crate::shell::format_rewind_execution_for_tui(&execution);
if let Some(error) = &session_record_error {
text.push_str("\nwarning: session recording failed: ");
text.push_str(error);
}
Some(Ok(TuiRewindWorkerResult::Execution {
execution,
text,
session_record_error,
}))
}
}
}
impl MissionControlApp {
pub(super) fn start_rewind_changes(&mut self, ui_state: &mut state::MissionControlState) {
let Some(session_id) = self.state.active_session_id().map(ToOwned::to_owned) else {
ui_state.status = "/changes requires an active persisted session".to_string();
return;
};
self.start_rewind_worker(
RewindRequestKind::Changes,
session_id,
RewindJob::Changes,
ui_state,
"loading rewindable changes…",
);
}
pub(super) fn start_rewind_command(
&mut self,
arg: Option<&str>,
ui_state: &mut state::MissionControlState,
) {
let Some(session_id) = self.state.active_session_id().map(ToOwned::to_owned) else {
ui_state.status = "/rewind requires an active persisted session".to_string();
return;
};
match crate::checkpoints::parse_rewind_target(arg) {
Ok(crate::checkpoints::ParsedRewindTarget::NeedsSelection) => {
self.start_rewind_worker(
RewindRequestKind::Plan,
session_id,
RewindJob::Selection,
ui_state,
"loading rewind targets…",
);
}
Ok(crate::checkpoints::ParsedRewindTarget::Target {
target_turn,
dry_run,
mode,
}) => {
self.start_rewind_worker(
RewindRequestKind::Plan,
session_id,
RewindJob::Plan {
target_turn,
dry_run,
mode,
},
ui_state,
"planning rewind…",
);
}
Err(error) => {
ui_state.status = format!(
"usage: /rewind [--dry-run] --to <turn> [--mode conversation|files|both]; {}",
bounded_rewind_diagnostic(error)
);
}
}
}
pub(super) fn start_rewind_execution(
&mut self,
plan: crate::checkpoints::RestorePlan,
mode: RewindMode,
ui_state: &mut state::MissionControlState,
) {
if self.active_run || self.worker.is_some() || self.pending_session_switch.is_some() {
ui_state.status = "cannot rewind while a prompt or session switch is running".into();
return;
}
let Some(session) = self.state.current_session.clone() else {
ui_state.status = "rewind requires an active persisted session".to_string();
return;
};
if session.id() != plan.session_id {
ui_state.status = "rewind target is no longer the active session".to_string();
return;
}
self.start_rewind_worker(
RewindRequestKind::Execute,
session.id().to_string(),
RewindJob::Execute {
plan,
session: Some(session),
mode,
},
ui_state,
"applying rewind…",
);
}
fn start_rewind_worker(
&mut self,
kind: RewindRequestKind,
session_id: String,
job: RewindJob,
ui_state: &mut state::MissionControlState,
status: &str,
) {
if self.pending_rewind.is_some() {
ui_state.status = "rewind request already in progress".to_string();
return;
}
self.next_rewind_request_id = self.next_rewind_request_id.saturating_add(1);
let request_id = self.next_rewind_request_id;
let generation = self.session_generation;
self.pending_rewind = Some(PendingRewind {
request_id,
session_id: session_id.clone(),
generation,
kind,
});
ui_state.status = status.to_string();
let paths = self.config.paths.clone();
let cwd = self.state.cwd.clone();
let sender = self.events.clone();
let cancel = Arc::new(AtomicBool::new(false));
let worker_cancel = Arc::clone(&cancel);
let worker_session_id = session_id.clone();
let handle = thread::Builder::new()
.name("magi-rewind".to_string())
.spawn(move || {
let Some(result) =
run_rewind_job(paths, worker_session_id.clone(), cwd, job, &worker_cancel)
else {
return RewindWorkerOutcome {
result: None,
delivery: RewindWorkerDelivery::Canceled,
};
};
let delivery = match send_tui_event(
&sender,
TuiEvent::RewindFinished {
request_id,
session_id: worker_session_id,
generation,
result: Box::new(result.clone()),
},
) {
Ok(()) => RewindWorkerDelivery::Delivered,
Err(_) => RewindWorkerDelivery::Failed,
};
RewindWorkerOutcome {
result: Some(result),
delivery,
}
});
match handle {
Ok(handle) => self.rewind_workers.push(RewindWorker {
request_id,
kind,
cancel,
handle,
}),
Err(error) => {
self.pending_rewind = None;
ui_state.status = format!(
"rewind failed: could not start worker: {}",
bounded_rewind_diagnostic(error)
);
}
}
}
pub(in crate::tui) fn handle_rewind_drain(
&mut self,
ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
) -> bool {
let mut changed = false;
for (request_id, session_id, generation, result) in &drain_result.rewind_finished {
let Some(pending) = self.pending_rewind.clone() else {
continue;
};
if pending.request_id != *request_id
|| pending.session_id != *session_id
|| pending.generation != *generation
{
continue;
}
self.pending_rewind = None;
if self.session_generation != *generation
|| self.state.active_session_id() != Some(session_id.as_str())
{
ui_state.status = "rewind result ignored; session changed".to_string();
changed = true;
continue;
}
if let Ok(worker_result) = result
&& rewind_worker_result_kind(worker_result) != pending.kind
&& !(pending.kind == RewindRequestKind::Plan
&& matches!(worker_result, TuiRewindWorkerResult::Changes { .. }))
{
ui_state.status = "rewind result ignored; request changed".to_string();
changed = true;
continue;
}
match result {
Ok(TuiRewindWorkerResult::Changes { text, plans }) => {
if plans.is_empty() {
ui_state.open_rewind_modal(text.clone(), None, true);
ui_state.status = if pending.kind == RewindRequestKind::Changes {
"showing rewindable changes"
} else {
"no prompts available for rewind"
}
.to_string();
} else {
let preview_texts =
plans.iter().map(|preview| preview.text.clone()).collect();
let plans = plans.iter().map(|preview| preview.plan.clone()).collect();
ui_state.open_rewind_picker_with_previews(plans, preview_texts);
ui_state.status = "select rewind target; press Enter to apply".to_string();
}
}
Ok(TuiRewindWorkerResult::Plan {
plan,
text,
dry_run,
mode,
}) => {
ui_state.open_rewind_modal(
text.clone(),
(!*dry_run).then_some(plan.clone()),
*dry_run,
);
if let Some(modal) = ui_state.modals.rewind_modal.as_mut() {
modal.mode = *mode;
modal.show_mode_choices = true;
}
ui_state.status = if *dry_run {
"showing rewind dry-run".to_string()
} else {
"previewing rewind; press Enter to apply".to_string()
};
}
Ok(TuiRewindWorkerResult::Execution {
execution,
text,
session_record_error,
}) => {
let mut status =
format!("rewind complete: {} operation(s)", execution.results.len());
if let Some(error) = session_record_error {
status.push_str("; session recording failed: ");
status.push_str(error);
}
ui_state.status = status;
ui_state.open_rewind_modal(text.clone(), None, true);
}
Ok(TuiRewindWorkerResult::Forked {
session_id,
session,
text,
}) => {
self.start_session_switch_worker_with_admission(
session_id.clone(),
Some(session.clone()),
ui_state,
);
ui_state.status = text.clone();
}
Err(error) => {
ui_state.status =
format!("rewind failed: {}", bounded_rewind_diagnostic(error));
}
}
changed = true;
}
changed
}
pub(in crate::tui) fn reap_rewind_workers(
&mut self,
ui_state: &mut state::MissionControlState,
) -> bool {
let mut changed = false;
let mut active = Vec::with_capacity(self.rewind_workers.len());
for worker in std::mem::take(&mut self.rewind_workers) {
if !worker.handle.is_finished() {
active.push(worker);
continue;
}
match worker.handle.join() {
Ok(outcome) => match outcome.delivery {
RewindWorkerDelivery::Delivered => {}
RewindWorkerDelivery::Canceled => {
if self
.pending_rewind
.as_ref()
.is_some_and(|p| p.request_id == worker.request_id)
{
self.pending_rewind = None;
ui_state.status = "rewind request cancelled".to_string();
changed = true;
}
}
RewindWorkerDelivery::Failed => {
if let Some(result) = outcome.result {
let (session_id, generation) = self
.pending_rewind
.as_ref()
.filter(|p| p.request_id == worker.request_id)
.map(|p| (p.session_id.clone(), p.generation))
.unwrap_or_default();
let mut drain = DrainResult::default();
drain.rewind_finished.push((
worker.request_id,
session_id,
generation,
result,
));
changed |= self.handle_rewind_drain(ui_state, &drain);
}
}
},
Err(_) => {
if self
.pending_rewind
.as_ref()
.is_some_and(|p| p.request_id == worker.request_id)
{
self.pending_rewind = None;
ui_state.status = "rewind worker panicked".to_string();
changed = true;
}
}
}
}
self.rewind_workers = active;
changed
}
pub(super) fn cancel_rewind_workers(&self) {
for worker in &self.rewind_workers {
worker.cancel.store(true, Ordering::SeqCst);
}
}
pub(in crate::tui) fn join_rewind_workers_for_cleanup(&mut self) -> Vec<String> {
let mut errors = Vec::new();
for worker in std::mem::take(&mut self.rewind_workers) {
if matches!(worker.kind, RewindRequestKind::Execute) {
if worker.handle.join().is_err() {
errors.push(format!("rewind worker {} panicked", worker.request_id));
}
continue;
}
let deadline = Instant::now() + WORKER_EXIT_JOIN_TIMEOUT;
let mut handle = Some(worker.handle);
while !handle.as_ref().is_some_and(JoinHandle::is_finished) {
let now = Instant::now();
if now >= deadline {
errors.push(format!(
"rewind worker {} did not stop within {}ms; detached",
worker.request_id,
WORKER_EXIT_JOIN_TIMEOUT.as_millis()
));
drop(handle.take());
break;
}
thread::sleep((deadline - now).min(Duration::from_millis(1)));
}
if let Some(handle) = handle
&& handle.join().is_err()
{
errors.push(format!("rewind worker {} panicked", worker.request_id));
}
}
self.pending_rewind = None;
errors
}
}