magi-code 0.96.2

Repository-aware CLI coding agent for terminal work
Documentation
use super::{Session, SessionEvent, SessionEventKind, SessionManager};

/// Uses the same durable user-input count as file checkpoints, including compaction.
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct RewindPrompt {
    pub(crate) turn: u64,
    pub(crate) text: String,
}

fn prompt_indices(events: &[SessionEvent]) -> Vec<(usize, u64)> {
    let mut turn = 0u64;
    let mut prompts = Vec::new();
    for (index, event) in events.iter().enumerate() {
        match event.kind() {
            Some(SessionEventKind::Compaction) => {
                if let Some(count) = event
                    .payload
                    .pointer("/aggregate/user_input_count")
                    .and_then(serde_json::Value::as_u64)
                {
                    turn = count;
                }
            }
            Some(SessionEventKind::UserInput) => {
                turn = turn.saturating_add(1);
                if event
                    .payload
                    .get("auto_recovery")
                    .and_then(serde_json::Value::as_bool)
                    != Some(true)
                {
                    prompts.push((index, turn));
                }
            }
            _ => {}
        }
    }
    prompts
}

fn prompt_text(event: &SessionEvent) -> anyhow::Result<&str> {
    event
        .payload
        .get("text")
        .and_then(serde_json::Value::as_str)
        .ok_or_else(|| anyhow::anyhow!("user prompt is missing text"))
}

impl Session {
    fn rewind_events(&self) -> anyhow::Result<Vec<SessionEvent>> {
        let append_lock = super::write::session_append_lock(self.path())?;
        let _guard = append_lock
            .lock()
            .map_err(|_| anyhow::anyhow!("session append lock was poisoned"))?;
        let _file_guard = crate::persistence::CrossProcessFileLock::acquire(self.path())?;
        let read =
            self.read_events_tolerant_bounded(super::MAX_METADATA_VISIT_LINES, 32 * 1024 * 1024)?;
        anyhow::ensure!(
            read.diagnostics.is_empty(),
            "cannot rewind unreadable session history"
        );
        anyhow::ensure!(
            read.events
                .iter()
                .all(|event| event.session_id == self.id()),
            "cannot rewind history containing foreign session records"
        );
        let (_, diagnostics) =
            super::latest_valid_compaction_checkpoint_for_replay(self.id(), &read.events);
        anyhow::ensure!(
            diagnostics.is_empty(),
            "cannot rewind malformed compaction history"
        );
        Ok(read.events)
    }

    pub(crate) fn rewind_prompts(&self) -> anyhow::Result<Vec<RewindPrompt>> {
        let events = self.rewind_events()?;
        prompt_indices(&events)
            .into_iter()
            .map(|(index, turn)| {
                Ok(RewindPrompt {
                    turn,
                    text: prompt_text(&events[index])?.to_owned(),
                })
            })
            .collect()
    }

    /// Fork strictly before the selected prompt; never rewrite the source session.
    /// Event-count compaction cutoffs remain valid because the prefix keeps its order.
    pub(crate) fn fork_before_prompt(
        &self,
        manager: &SessionManager,
        target: u64,
    ) -> anyhow::Result<Session> {
        let mut events = self.rewind_events()?;
        let index = prompt_indices(&events)
            .into_iter()
            .find_map(|(index, turn)| (turn == target).then_some(index))
            .ok_or_else(|| anyhow::anyhow!("prompt target is not in retained session history"))?;
        let prompt = prompt_text(&events[index])?.to_owned();
        let cwd = events[index].cwd.clone();
        events.truncate(index);
        let fork = manager.create()?.admit_standalone_writer()?;
        for event in &mut events {
            event.session_id = fork.id().to_owned();
        }
        events.push(SessionEvent::new_kind(
            SessionEventKind::Diagnostic,
            fork.id().to_owned(),
            cwd,
            serde_json::json!({"rewind_source_session": self.id(), "before_user_turn": target, "rewind_prompt": prompt}),
        ));
        fork.append_owned_batch(events).map_err(|error| {
            anyhow::anyhow!("could not persist rewind fork {}: {}", fork.id(), error)
        })?;
        Ok(fork)
    }
}