magi-code 0.77.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::*;
use std::io::Write;

fn append_completed_turn(session: &Session, cwd: &Path, user: &str, assistant: &str) {
    for (kind, payload) in [
        (SessionEventKind::UserInput, json!({"text": user})),
        (
            SessionEventKind::AssistantOutput,
            json!({"text": assistant}),
        ),
        (SessionEventKind::TurnStatus, json!({"status": "completed"})),
    ] {
        session
            .append(&SessionEvent::new_kind(
                kind,
                session.id().to_string(),
                cwd.to_path_buf(),
                payload,
            ))
            .unwrap();
    }
}

fn replay_text(replay: &ConversationReplay) -> String {
    format!("{:?}", replay.items)
}

fn write_private(path: &Path, bytes: &[u8]) {
    fs::write(path, bytes).unwrap();
    #[cfg(unix)]
    {
        use std::os::unix::fs::PermissionsExt;
        fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
    }
}

#[test]
fn unchanged_snapshot_reuses_replay_without_second_scan() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "question", "answer");
    let mut cache = ConversationReplayCache::default();

    let first = cache.replay(Some(&session)).unwrap();
    let second = cache.replay(Some(&session)).unwrap();

    assert_eq!(first, second);
    assert_eq!(cache.metrics().full_scans, 1);
    assert_eq!(cache.metrics().replay_builds, 1);
    assert_eq!(cache.metrics().cache_hits, 1);
    assert!(cache.metrics().bytes_read > 0);
    assert_eq!(cache.metrics().lines_parsed, 3);
    assert_eq!(cache.metrics().events_parsed, 3);
}

#[test]
fn durable_append_invalidates_snapshot_and_is_visible() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "first", "one");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();

    append_completed_turn(&session, temp.path(), "second", "two");
    let replay = cache.replay(Some(&session)).unwrap();

    assert!(replay_text(&replay).contains("second"));
    assert!(replay_text(&replay).contains("two"));
    assert_eq!(cache.metrics().full_scans, 2);
    assert_eq!(cache.metrics().invalidations, 1);
}

#[test]
fn external_append_invalidates_snapshot_and_is_visible() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "first", "one");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();

    let event = SessionEvent::new_kind(
        SessionEventKind::UserInput,
        session.id().to_string(),
        temp.path().to_path_buf(),
        json!({"text": "external"}),
    );
    let mut file = fs::OpenOptions::new()
        .append(true)
        .open(session.path())
        .unwrap();
    serde_json::to_writer(&mut file, &event).unwrap();
    file.write_all(b"\n").unwrap();
    file.sync_all().unwrap();

    let replay = cache.replay(Some(&session)).unwrap();
    assert!(replay_text(&replay).contains("external"));
    assert_eq!(cache.metrics().full_scans, 2);
}

#[cfg(unix)]
#[test]
fn same_size_rewrite_with_restored_mtime_cannot_reuse_stale_replay() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "alpha", "reply");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();

    let old_modified = fs::metadata(session.path()).unwrap().modified().unwrap();
    let original = fs::read_to_string(session.path()).unwrap();
    let replacement = original.replacen("alpha", "bravo", 1);
    assert_eq!(original.len(), replacement.len());
    fs::write(session.path(), replacement).unwrap();
    fs::File::options()
        .write(true)
        .open(session.path())
        .unwrap()
        .set_modified(old_modified)
        .unwrap();

    let replay = cache.replay(Some(&session)).unwrap();
    assert!(replay_text(&replay).contains("bravo"));
    assert!(!replay_text(&replay).contains("alpha"));
    assert_eq!(cache.metrics().full_scans, 2);
}

#[test]
fn partial_trailing_record_completed_between_reads_is_revalidated() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    let event = SessionEvent::new_kind(
        SessionEventKind::UserInput,
        session.id().to_string(),
        temp.path().to_path_buf(),
        json!({"text": "completed later"}),
    );
    let encoded = serde_json::to_vec(&event).unwrap();
    let split = encoded.len() / 2;
    write_private(session.path(), &encoded[..split]);
    let mut cache = ConversationReplayCache::default();

    let partial = cache.replay(Some(&session)).unwrap();
    assert_eq!(partial.session_read_diagnostics.len(), 1);
    let mut file = fs::OpenOptions::new()
        .append(true)
        .open(session.path())
        .unwrap();
    file.write_all(&encoded[split..]).unwrap();
    file.write_all(b"\n").unwrap();
    file.sync_all().unwrap();

    let completed = cache.replay(Some(&session)).unwrap();
    assert!(completed.session_read_diagnostics.is_empty());
    assert!(replay_text(&completed).contains("completed later"));
}

#[test]
fn concurrent_append_is_visible_on_the_next_replay() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "first", "reply");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();
    let writer = session.clone();
    let cwd = temp.path().to_path_buf();

    std::thread::spawn(move || append_completed_turn(&writer, &cwd, "concurrent", "visible"))
        .join()
        .unwrap();

    let replay = cache.replay(Some(&session)).unwrap();
    assert!(replay_text(&replay).contains("concurrent"));
    assert!(replay_text(&replay).contains("visible"));
}

#[test]
fn malformed_line_diagnostics_are_reused_once_in_stable_order() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "valid", "reply");
    let mut file = fs::OpenOptions::new()
        .append(true)
        .open(session.path())
        .unwrap();
    file.write_all(b"{malformed}\n").unwrap();
    file.sync_all().unwrap();
    let mut cache = ConversationReplayCache::default();

    let first = cache.replay(Some(&session)).unwrap();
    let second = cache.replay(Some(&session)).unwrap();

    assert_eq!(
        first.session_read_diagnostics,
        second.session_read_diagnostics
    );
    assert_eq!(second.session_read_diagnostics.len(), 1);
    assert_eq!(cache.metrics().full_scans, 1);
    assert_eq!(cache.metrics().cache_hits, 1);
}

#[test]
fn physical_rotation_invalidates_cached_replay() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "old", "old reply");
    append_completed_turn(&session, temp.path(), "recent", "recent reply");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();

    crate::sessions::record_session_compaction(
        &session,
        temp.path(),
        "summary of old turn",
        "provider",
        "model",
        3,
    )
    .unwrap();

    let replay = cache.replay(Some(&session)).unwrap();
    let text = replay_text(&replay);
    assert!(text.contains("summary of old turn"));
    assert!(text.contains("recent"));
    assert!(!text.contains("old reply"));
    assert_eq!(cache.metrics().full_scans, 2);
}

#[test]
fn truncation_and_replacement_invalidate_replay() {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    append_completed_turn(&session, temp.path(), "old", "reply");
    let mut cache = ConversationReplayCache::default();
    cache.replay(Some(&session)).unwrap();

    fs::write(session.path(), []).unwrap();
    assert!(cache.replay(Some(&session)).unwrap().items.is_empty());

    let event = SessionEvent::new_kind(
        SessionEventKind::UserInput,
        session.id().to_string(),
        temp.path().to_path_buf(),
        json!({"text": "replacement"}),
    );
    let replacement = session.path().with_extension("replacement");
    let mut bytes = serde_json::to_vec(&event).unwrap();
    bytes.push(b'\n');
    write_private(&replacement, &bytes);
    fs::rename(&replacement, session.path()).unwrap();
    let replay = cache.replay(Some(&session)).unwrap();
    assert!(replay_text(&replay).contains("replacement"));
    assert_eq!(cache.metrics().full_scans, 3);
}

fn add_metrics(total: &mut ReplayCacheMetrics, next: ReplayCacheMetrics) {
    total.full_scans += next.full_scans;
    total.bytes_read += next.bytes_read;
    total.lines_parsed += next.lines_parsed;
    total.events_parsed += next.events_parsed;
    total.replay_builds += next.replay_builds;
    total.events_before_cutoff += next.events_before_cutoff;
    total.events_after_cutoff += next.events_after_cutoff;
    total.cache_hits += next.cache_hits;
    total.invalidations += next.invalidations;
}

fn replay_benchmark_sample() -> (
    std::time::Duration,
    std::time::Duration,
    ReplayCacheMetrics,
    ReplayCacheMetrics,
) {
    let temp = TempDir::new().unwrap();
    let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
        .create()
        .unwrap();
    let mut baseline_elapsed = std::time::Duration::ZERO;
    let mut optimized_elapsed = std::time::Duration::ZERO;
    let mut baseline_metrics = ReplayCacheMetrics::default();
    let mut optimized = ConversationReplayCache::default();
    const PROMPTS: usize = 40;

    for index in 0..PROMPTS {
        append_completed_turn(
            &session,
            temp.path(),
            &format!("question {index} {}", "x".repeat(2_000)),
            &format!("answer {index} {}", "y".repeat(2_000)),
        );

        let start = std::time::Instant::now();
        for _ in 0..2 {
            let mut replay = ConversationReplayCache::default();
            replay.replay(Some(&session)).unwrap();
            add_metrics(&mut baseline_metrics, replay.metrics());
        }
        baseline_elapsed += start.elapsed();

        let start = std::time::Instant::now();
        optimized.replay(Some(&session)).unwrap();
        optimized.replay(Some(&session)).unwrap();
        optimized_elapsed += start.elapsed();
    }
    (
        baseline_elapsed,
        optimized_elapsed,
        baseline_metrics,
        optimized.metrics(),
    )
}

fn percentile(samples: &mut [std::time::Duration], percentile: usize) -> std::time::Duration {
    samples.sort_unstable();
    let index = (samples.len() - 1) * percentile / 100;
    samples[index]
}

#[test]
#[ignore = "release-mode replay benchmark; run with --release --ignored --nocapture"]
fn replay_snapshot_benchmark_growing_sessions() {
    const WARMUP_SAMPLES: usize = 1;
    const SAMPLES: usize = 9;
    const PROMPTS: u64 = 40;
    for _ in 0..WARMUP_SAMPLES {
        let _ = replay_benchmark_sample();
    }
    let mut baseline_times = Vec::with_capacity(SAMPLES);
    let mut optimized_times = Vec::with_capacity(SAMPLES);
    let mut baseline_metrics = ReplayCacheMetrics::default();
    let mut optimized_metrics = ReplayCacheMetrics::default();
    for _ in 0..SAMPLES {
        let (baseline, optimized, baseline_sample, optimized_sample) = replay_benchmark_sample();
        baseline_times.push(baseline);
        optimized_times.push(optimized);
        baseline_metrics = baseline_sample;
        optimized_metrics = optimized_sample;
    }
    let baseline_median = percentile(&mut baseline_times.clone(), 50);
    let baseline_p95 = percentile(&mut baseline_times, 95);
    let optimized_median = percentile(&mut optimized_times.clone(), 50);
    let optimized_p95 = percentile(&mut optimized_times, 95);

    eprintln!(
        "session_replay_benchmark fixture_seed=sequential-0.70 warmup={WARMUP_SAMPLES} samples={SAMPLES} prompts={PROMPTS} baseline_median_ms={} baseline_p95_ms={} optimized_median_ms={} optimized_p95_ms={} baseline_scans={} optimized_scans={} baseline_bytes={} optimized_bytes={} baseline_lines={} optimized_lines={} baseline_events={} optimized_events={} baseline_replay_builds={} optimized_replay_builds={} optimized_cache_hits={}",
        baseline_median.as_millis(),
        baseline_p95.as_millis(),
        optimized_median.as_millis(),
        optimized_p95.as_millis(),
        baseline_metrics.full_scans,
        optimized_metrics.full_scans,
        baseline_metrics.bytes_read,
        optimized_metrics.bytes_read,
        baseline_metrics.lines_parsed,
        optimized_metrics.lines_parsed,
        baseline_metrics.events_parsed,
        optimized_metrics.events_parsed,
        baseline_metrics.replay_builds,
        optimized_metrics.replay_builds,
        optimized_metrics.cache_hits,
    );
    assert_eq!(baseline_metrics.full_scans, PROMPTS * 2);
    assert_eq!(optimized_metrics.full_scans, PROMPTS);
    assert_eq!(optimized_metrics.replay_builds, PROMPTS);
    assert_eq!(optimized_metrics.cache_hits, PROMPTS);
    // Chosen before implementation: median wall time must improve by at least 25%.
    assert!(optimized_median * 4 <= baseline_median * 3);
}