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);
assert!(optimized_median * 4 <= baseline_median * 3);
}