use super::{
capture::{CHECKPOINT_SCHEMA_VERSION, LedgerRead, SnapshotBlobRef, SnapshotLedgerRecord},
sanitized_diagnostic,
};
use crate::{
config::McPaths,
persistence::{CrossProcessFileLock, atomic_write, sync_parent_dir},
};
use sha2::{Digest, Sha256};
use std::{
error::Error,
fmt, fs,
io::{ErrorKind, Read, Seek, SeekFrom, Write},
path::{Path, PathBuf},
};
#[derive(Debug)]
pub(crate) enum CheckpointAppendFailure {
OutcomeUncertain {
path: PathBuf,
source: anyhow::Error,
},
CommittedButUndurable {
path: PathBuf,
source: anyhow::Error,
},
}
impl CheckpointAppendFailure {
fn uncertain(path: &Path, source: anyhow::Error) -> anyhow::Error {
anyhow::Error::new(Self::OutcomeUncertain {
path: path.to_path_buf(),
source,
})
}
fn committed_but_undurable(path: &Path, source: anyhow::Error) -> anyhow::Error {
anyhow::Error::new(Self::CommittedButUndurable {
path: path.to_path_buf(),
source,
})
}
}
impl fmt::Display for CheckpointAppendFailure {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
let (path, outcome) = match self {
Self::OutcomeUncertain { path, .. } => {
(path, "append outcome is uncertain; do not retry blindly")
}
Self::CommittedButUndurable { path, .. } => (
path,
"record is visible but final synchronization failed; do not retry blindly",
),
};
write!(
formatter,
"checkpoint ledger append failed for {}: {outcome}",
path.display()
)
}
}
impl Error for CheckpointAppendFailure {
fn source(&self) -> Option<&(dyn Error + 'static)> {
let source = match self {
Self::OutcomeUncertain { source, .. } | Self::CommittedButUndurable { source, .. } => {
source
}
};
Some(source.as_ref())
}
}
#[derive(Debug, Default)]
struct LedgerIoMetrics {
read_bytes: u64,
write_bytes: u64,
serialized_record_bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct CheckpointStore {
pub(super) root: PathBuf,
}
impl CheckpointStore {
pub(crate) fn new(root: PathBuf) -> Self {
Self { root }
}
pub(crate) fn from_paths(paths: &McPaths) -> Self {
Self::new(paths.checkpoints.clone())
}
pub(crate) fn copy_before_turn(
&self,
source: &str,
target: &str,
turn: u64,
) -> anyhow::Result<()> {
crate::sessions::validate_session_id(source.to_owned())?;
crate::sessions::validate_session_id(target.to_owned())?;
let read = self.read_records(source);
anyhow::ensure!(
read.diagnostics.is_empty(),
"cannot copy unreadable file checkpoints"
);
for mut record in read
.records
.into_iter()
.filter(|record| record.event.user_turn < turn)
{
record.event.session_id = target.to_owned();
self.append_record(&record)?;
}
Ok(())
}
pub(super) fn blobs_dir(&self) -> PathBuf {
self.root.join("blobs")
}
pub(super) fn ledgers_dir(&self) -> PathBuf {
self.root.join("ledgers")
}
pub(super) fn ledger_path(&self, session_id: &str) -> PathBuf {
self.ledgers_dir().join(format!("{session_id}.jsonl"))
}
pub(super) fn blob_path(&self, sha256: &str) -> PathBuf {
self.blobs_dir().join(sha256)
}
fn mutation_lock_target(&self) -> PathBuf {
self.root.with_file_name("checkpoint-store")
}
pub(super) fn lock_mutations(&self) -> anyhow::Result<CrossProcessFileLock> {
CrossProcessFileLock::acquire(&self.mutation_lock_target())
}
pub(super) fn write_blob(&self, bytes: &[u8]) -> anyhow::Result<SnapshotBlobRef> {
let sha256 = sha256_hex(bytes);
let path = self.blob_path(&sha256);
if path.exists() {
return Ok(SnapshotBlobRef {
sha256,
bytes: bytes.len() as u64,
});
}
atomic_write(&path, bytes)?;
Ok(SnapshotBlobRef {
sha256,
bytes: bytes.len() as u64,
})
}
pub(super) fn read_blob(&self, blob: &SnapshotBlobRef) -> anyhow::Result<Vec<u8>> {
let bytes = fs::read(self.blob_path(&blob.sha256))?;
if sha256_hex(&bytes) != blob.sha256 {
anyhow::bail!("checkpoint blob hash mismatch");
}
Ok(bytes)
}
pub(super) fn append_record(&self, record: &SnapshotLedgerRecord) -> anyhow::Result<()> {
self.append_record_with(
record,
None,
|file| file.sync_all().map_err(anyhow::Error::from),
sync_parent_dir,
)
}
fn append_record_with(
&self,
record: &SnapshotLedgerRecord,
mut metrics: Option<&mut LedgerIoMetrics>,
sync_file: impl FnOnce(&fs::File) -> anyhow::Result<()>,
sync_parent: impl FnOnce(&Path) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
crate::sessions::validate_session_id(record.event.session_id.clone())?;
let ledger = self.ledger_path(&record.event.session_id);
if let Some(parent) = ledger.parent() {
fs::create_dir_all(parent)?;
}
let mut line = serde_json::to_vec(record)?;
line.push(b'\n');
if let Some(metrics) = metrics.as_deref_mut() {
metrics.serialized_record_bytes += line.len() as u64;
}
let _lock = CrossProcessFileLock::acquire(&ledger)?;
let (mut file, created) = open_ledger_for_append(&ledger)?;
let pre_append_len = file.metadata()?.len();
let needs_separator = if pre_append_len == 0 {
false
} else {
file.seek(SeekFrom::End(-1))?;
let mut tail = [0_u8; 1];
file.read_exact(&mut tail)?;
if let Some(metrics) = metrics.as_deref_mut() {
metrics.read_bytes += 1;
}
tail[0] != b'\n'
};
if needs_separator {
line.insert(0, b'\n');
}
if let Some(metrics) = metrics {
metrics.write_bytes += line.len() as u64;
}
if let Err(write_error) = file.write_all(&line) {
let rollback = file.set_len(pre_append_len).and_then(|()| file.sync_all());
return match rollback {
Ok(()) => Err(anyhow::Error::new(write_error).context(
"checkpoint ledger append failed; partial data was durably rolled back",
)),
Err(rollback_error) => Err(CheckpointAppendFailure::uncertain(
&ledger,
anyhow::anyhow!(
"append write failed ({write_error}); rollback failed ({rollback_error})"
),
)),
};
}
file.flush().map_err(|error| {
CheckpointAppendFailure::committed_but_undurable(&ledger, error.into())
})?;
sync_file(&file)
.map_err(|error| CheckpointAppendFailure::committed_but_undurable(&ledger, error))?;
if created {
let parent = ledger
.parent()
.ok_or_else(|| anyhow::anyhow!("checkpoint ledger has no parent"))?;
sync_parent(parent).map_err(|error| {
CheckpointAppendFailure::committed_but_undurable(&ledger, error)
})?;
}
Ok(())
}
pub(crate) fn read_records(&self, session_id: &str) -> LedgerRead {
let ledger = self.ledger_path(session_id);
let mut read = LedgerRead::default();
let bytes = match fs::read(&ledger) {
Ok(bytes) => bytes,
Err(error) if error.kind() == ErrorKind::NotFound => return read,
Err(error) => {
read.diagnostics.push(sanitized_diagnostic(format!(
"checkpoint ledger read failed: {error}"
)));
return read;
}
};
let final_segment = bytes.split(|byte| *byte == b'\n').count().saturating_sub(1);
let has_terminated_tail = bytes.ends_with(b"\n");
for (index, line) in bytes.split(|byte| *byte == b'\n').enumerate() {
if line.iter().all(u8::is_ascii_whitespace) {
continue;
}
match serde_json::from_slice::<SnapshotLedgerRecord>(line) {
Ok(record) if record.schema_version == CHECKPOINT_SCHEMA_VERSION => {
read.records.push(record);
}
Ok(_) => read.diagnostics.push(format!(
"ignored checkpoint ledger line {}: unsupported_schema_version",
index + 1
)),
Err(_) if index == final_segment && !has_terminated_tail => {
read.diagnostics.push(format!(
"ignored checkpoint ledger line {}: incomplete_tail",
index + 1
))
}
Err(_) => read.diagnostics.push(format!(
"ignored checkpoint ledger line {}: malformed_json",
index + 1
)),
}
}
read
}
pub(crate) fn prune_session(&self, session_id: &str) -> anyhow::Result<()> {
if !self.validate_deletion_directories()? {
return Ok(());
}
let _mutation_lock = self.lock_mutations()?;
if !self.validate_deletion_directories()? {
return Ok(());
}
let ledger = self.ledger_path(session_id);
if ledger.exists() {
fs::remove_file(&ledger)?;
}
self.prune_unreferenced_blobs_locked()
}
pub(crate) fn prune_unreferenced_blobs(&self) -> anyhow::Result<()> {
if !self.validate_deletion_directories()? {
return Ok(());
}
let _mutation_lock = self.lock_mutations()?;
if !self.validate_deletion_directories()? {
return Ok(());
}
self.prune_unreferenced_blobs_locked()
}
fn validate_deletion_directories(&self) -> anyhow::Result<bool> {
if !validate_owned_directory(&self.root)? {
return Ok(false);
}
for path in [self.ledgers_dir(), self.blobs_dir()] {
match fs::symlink_metadata(&path) {
Ok(_) => {
if !validate_owned_directory(&path)? {
anyhow::bail!("unsafe checkpoint directory: {}", path.display());
}
}
Err(error) if error.kind() == ErrorKind::NotFound => {}
Err(error) => return Err(error.into()),
}
}
Ok(true)
}
fn prune_unreferenced_blobs_locked(&self) -> anyhow::Result<()> {
let mut referenced = std::collections::BTreeSet::new();
let ledgers = self.ledgers_dir();
if ledgers.exists() {
for entry in fs::read_dir(&ledgers)? {
let entry = entry?;
let file_type = entry.file_type()?;
if file_type.is_symlink() {
anyhow::bail!("unsafe checkpoint ledger entry");
}
if !file_type.is_file() {
continue;
}
let file_name = entry
.file_name()
.into_string()
.map_err(|_| anyhow::anyhow!("checkpoint ledger name is not UTF-8"))?;
let Some(session_id) = file_name.strip_suffix(".jsonl") else {
continue;
};
crate::sessions::validate_session_id(session_id.to_string())?;
let read = self.read_records(session_id);
if !read.diagnostics.is_empty() {
anyhow::bail!("checkpoint ledger could not be read completely");
}
for record in read.records {
if let Some(blob) = record.event.pre {
referenced.insert(blob.sha256);
}
if let Some(blob) = record.event.post {
referenced.insert(blob.sha256);
}
}
}
}
let blobs = self.blobs_dir();
if !blobs.exists() {
return Ok(());
}
for entry in fs::read_dir(&blobs)? {
let entry = entry?;
if !entry.file_type()?.is_file() {
continue;
}
let name = entry.file_name().to_string_lossy().into_owned();
if !referenced.contains(&name) {
fs::remove_file(entry.path())?;
}
}
Ok(())
}
}
fn open_ledger_for_append(path: &Path) -> anyhow::Result<(fs::File, bool)> {
match fs::symlink_metadata(path) {
Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_file() => {
anyhow::bail!("unsafe checkpoint ledger entry")
}
Ok(_) => {
let file = fs::OpenOptions::new().read(true).append(true).open(path)?;
ensure_open_ledger_matches_path(&file, path)?;
Ok((file, false))
}
Err(error) if error.kind() == ErrorKind::NotFound => {
let file = fs::OpenOptions::new()
.read(true)
.append(true)
.create_new(true)
.open(path)?;
ensure_open_ledger_matches_path(&file, path)?;
Ok((file, true))
}
Err(error) => Err(error.into()),
}
}
fn ensure_open_ledger_matches_path(file: &fs::File, path: &Path) -> anyhow::Result<()> {
let path_metadata = fs::symlink_metadata(path)?;
let file_metadata = file.metadata()?;
if path_metadata.file_type().is_symlink()
|| !path_metadata.is_file()
|| !file_metadata.is_file()
{
anyhow::bail!("unsafe checkpoint ledger entry");
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if path_metadata.dev() != file_metadata.dev()
|| path_metadata.ino() != file_metadata.ino()
|| file_metadata.uid() != unsafe { libc::geteuid() }
{
anyhow::bail!("unsafe checkpoint ledger entry");
}
}
Ok(())
}
fn validate_owned_directory(path: &std::path::Path) -> anyhow::Result<bool> {
let metadata = match fs::symlink_metadata(path) {
Ok(metadata) => metadata,
Err(error) if error.kind() == ErrorKind::NotFound => return Ok(false),
Err(error) => return Err(error.into()),
};
if metadata.file_type().is_symlink() || !metadata.file_type().is_dir() {
anyhow::bail!("unsafe checkpoint directory: {}", path.display());
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if metadata.uid() != unsafe { libc::geteuid() } {
anyhow::bail!("checkpoint directory is not owned by the current user");
}
}
Ok(true)
}
pub(crate) fn sha256_hex(bytes: &[u8]) -> String {
let digest = Sha256::digest(bytes);
crate::hex::lower_hex(digest)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::checkpoints::{
SnapshotContext, SnapshotExclusionReason, SnapshotTool,
capture::capture_file_snapshot_with_hook,
capture::{FileSnapshotEvent, SnapshotCaptureStatus, SnapshotLedgerRecord},
};
use std::{sync::mpsc, thread, time::Duration};
use tempfile::TempDir;
fn store(temp: &TempDir) -> CheckpointStore {
CheckpointStore::new(temp.path().join("mc/checkpoints"))
}
fn context(temp: &TempDir) -> SnapshotContext {
let paths = McPaths::from_root(temp.path().join("mc"));
SnapshotContext {
store: CheckpointStore::from_paths(&paths),
paths,
session_id: "session".to_string(),
user_turn: 1,
}
}
fn test_record(session_id: &str, cwd: &Path, turn: u64) -> SnapshotLedgerRecord {
SnapshotLedgerRecord::new(FileSnapshotEvent {
session_id: session_id.to_string(),
user_turn: turn,
tool: SnapshotTool::Edit,
cwd: cwd.to_path_buf(),
relative_path: PathBuf::from("file.txt"),
pre: None,
post: None,
status: SnapshotCaptureStatus::Excluded {
reason: SnapshotExclusionReason::NotFile,
},
})
}
#[test]
fn blob_store_round_trips_by_hash() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let blob = store.write_blob(b"hello").unwrap();
assert_eq!(store.read_blob(&blob).unwrap(), b"hello");
assert_eq!(blob.sha256, sha256_hex(b"hello"));
}
#[test]
fn blob_write_is_idempotent() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let first = store.write_blob(b"same").unwrap();
let second = store.write_blob(b"same").unwrap();
assert_eq!(first, second);
}
#[test]
fn checkpoint_gc_waits_for_blob_publication() {
let temp = TempDir::new().unwrap();
let context = context(&temp);
let old_record = SnapshotLedgerRecord::new(FileSnapshotEvent {
session_id: "old-session".to_string(),
user_turn: 1,
tool: SnapshotTool::Edit,
cwd: temp.path().to_path_buf(),
relative_path: PathBuf::from("old.txt"),
pre: None,
post: None,
status: SnapshotCaptureStatus::Excluded {
reason: SnapshotExclusionReason::NotFile,
},
});
context.store.append_record(&old_record).unwrap();
let target = temp.path().join("new.txt");
fs::write(&target, "new").unwrap();
let (published_tx, published_rx) = mpsc::sync_channel(0);
let (release_tx, release_rx) = mpsc::sync_channel(0);
let capture_context = context.clone();
let capture_root = temp.path().to_path_buf();
let capture_target = target.clone();
let capture = thread::spawn(move || {
capture_file_snapshot_with_hook(
&capture_context,
&capture_root,
SnapshotTool::WriteFile,
&capture_target,
None,
Some(b"new"),
|| {
published_tx.send(()).unwrap();
release_rx.recv().unwrap();
},
)
});
published_rx.recv().unwrap();
let prune_store = context.store.clone();
let (started_tx, started_rx) = mpsc::sync_channel(0);
let (pruned_tx, pruned_rx) = mpsc::sync_channel(0);
let prune = thread::spawn(move || {
started_tx.send(()).unwrap();
pruned_tx
.send(prune_store.prune_session("old-session"))
.unwrap();
});
started_rx.recv().unwrap();
assert!(matches!(
pruned_rx.recv_timeout(Duration::from_millis(100)),
Err(mpsc::RecvTimeoutError::Timeout)
));
release_tx.send(()).unwrap();
capture.join().unwrap().unwrap();
pruned_rx.recv().unwrap().unwrap();
prune.join().unwrap();
let record = context.store.read_records("session").records.pop().unwrap();
let post = record.event.post.unwrap();
assert_eq!(context.store.read_blob(&post).unwrap(), b"new");
let orphan = context.store.blobs_dir().join("orphan-after-failed-gc");
fs::write(&orphan, "orphan").unwrap();
context.store.prune_unreferenced_blobs().unwrap();
assert!(!orphan.exists());
}
#[cfg(unix)]
#[test]
fn checkpoint_gc_rejects_symlinked_blob_directory() {
use std::os::unix::fs::symlink;
let temp = TempDir::new().unwrap();
let store = store(&temp);
fs::create_dir_all(&store.root).unwrap();
let outside = temp.path().join("outside");
fs::create_dir_all(&outside).unwrap();
let retained = outside.join("retained");
fs::write(&retained, "keep").unwrap();
symlink(&outside, store.blobs_dir()).unwrap();
assert!(store.prune_unreferenced_blobs().is_err());
assert_eq!(fs::read_to_string(retained).unwrap(), "keep");
}
#[test]
fn retention_safety_preserves_blobs_for_malformed_retained_ledger() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
fs::create_dir_all(store.blobs_dir()).unwrap();
fs::create_dir_all(store.ledgers_dir()).unwrap();
let retained = store.blobs_dir().join("retained");
fs::write(&retained, "checkpoint").unwrap();
fs::write(store.ledger_path("live-session"), "{malformed\n").unwrap();
assert!(store.prune_unreferenced_blobs().is_err());
assert_eq!(fs::read_to_string(retained).unwrap(), "checkpoint");
}
#[test]
fn ledger_skips_malformed_lines() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let ledger = store.ledger_path("session");
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(&ledger, "not json\n").unwrap();
let read = store.read_records("session");
assert!(read.records.is_empty());
assert_eq!(read.diagnostics.len(), 1);
}
#[test]
fn ledger_append_recovers_after_truncated_line() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let ledger = store.ledger_path("session");
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(&ledger, "truncated").unwrap();
let event = FileSnapshotEvent {
session_id: "session".to_string(),
user_turn: 1,
tool: SnapshotTool::WriteFile,
cwd: temp.path().to_path_buf(),
relative_path: PathBuf::from("file.txt"),
pre: None,
post: None,
status: SnapshotCaptureStatus::Excluded {
reason: SnapshotExclusionReason::SecretPath,
},
};
store
.append_record(&SnapshotLedgerRecord::new(event))
.unwrap();
let read = store.read_records("session");
assert_eq!(read.records.len(), 1);
assert_eq!(read.diagnostics.len(), 1);
}
#[test]
fn ledger_reader_preserves_valid_history_before_incomplete_tail() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
store
.append_record(&test_record("session", temp.path(), 1))
.unwrap();
let ledger = store.ledger_path("session");
let mut file = fs::OpenOptions::new().append(true).open(&ledger).unwrap();
file.write_all(b"{\"schema_version\":1").unwrap();
file.sync_all().unwrap();
let read = store.read_records("session");
assert_eq!(read.records.len(), 1);
assert_eq!(
read.diagnostics,
vec!["ignored checkpoint ledger line 2: incomplete_tail"]
);
}
#[test]
fn ledger_reader_surfaces_malformed_committed_middle_record() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let first = serde_json::to_vec(&test_record("session", temp.path(), 1)).unwrap();
let second = serde_json::to_vec(&test_record("session", temp.path(), 2)).unwrap();
let mut ledger_bytes = first;
ledger_bytes.extend_from_slice(b"\n{malformed}\n");
ledger_bytes.extend_from_slice(&second);
ledger_bytes.push(b'\n');
let ledger = store.ledger_path("session");
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(ledger, ledger_bytes).unwrap();
let read = store.read_records("session");
assert_eq!(read.records.len(), 2);
assert_eq!(
read.diagnostics,
vec!["ignored checkpoint ledger line 2: malformed_json"]
);
}
#[test]
fn complete_unterminated_record_is_read_and_separated_from_next_append() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let ledger = store.ledger_path("session");
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(
&ledger,
serde_json::to_vec(&test_record("session", temp.path(), 1)).unwrap(),
)
.unwrap();
store
.append_record(&test_record("session", temp.path(), 2))
.unwrap();
let read = store.read_records("session");
assert_eq!(read.records.len(), 2);
assert!(read.diagnostics.is_empty());
}
#[test]
fn visible_record_with_failed_sync_has_distinguishable_outcome() {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let record = test_record("session", temp.path(), 1);
let error = store
.append_record_with(
&record,
None,
|_| Err(anyhow::anyhow!("injected file sync failure")),
|_| Ok(()),
)
.unwrap_err();
let failure = error
.downcast_ref::<CheckpointAppendFailure>()
.expect("typed checkpoint append failure");
assert!(matches!(
failure,
CheckpointAppendFailure::CommittedButUndurable { .. }
));
assert_eq!(store.read_records("session").records, vec![record]);
}
#[test]
fn checkpoint_append_logical_io_scales_linearly() {
fn measure(count: usize) -> (LedgerIoMetrics, u64, Duration) {
let temp = TempDir::new().unwrap();
let store = store(&temp);
let record = test_record("session", temp.path(), 1);
let mut metrics = LedgerIoMetrics::default();
let started = std::time::Instant::now();
for _ in 0..count {
store
.append_record_with(&record, Some(&mut metrics), |_| Ok(()), |_| Ok(()))
.unwrap();
}
let final_bytes = fs::metadata(store.ledger_path("session")).unwrap().len();
(metrics, final_bytes, started.elapsed())
}
let (at_1k, final_1k, elapsed_1k) = measure(1_000);
let (at_2k, final_2k, elapsed_2k) = measure(2_000);
let logical_1k = at_1k.read_bytes + at_1k.write_bytes;
let logical_2k = at_2k.read_bytes + at_2k.write_bytes;
let ratio = logical_2k as f64 / logical_1k as f64;
assert_eq!(at_1k.serialized_record_bytes, final_1k);
assert_eq!(at_2k.serialized_record_bytes, final_2k);
assert!(ratio <= 2.5, "logical-byte ratio was {ratio}");
eprintln!(
"checkpoint ledger logical I/O: 1k read={} write={} final={} elapsed={elapsed_1k:?}; 2k read={} write={} final={} elapsed={elapsed_2k:?}; ratio={ratio:.3}",
at_1k.read_bytes,
at_1k.write_bytes,
final_1k,
at_2k.read_bytes,
at_2k.write_bytes,
final_2k,
);
}
#[test]
fn ledger_process_append_helper() {
let Some(root) = std::env::var_os("MAGI_CHECKPOINT_APPEND_HELPER_ROOT") else {
return;
};
let writer = std::env::var("MAGI_CHECKPOINT_APPEND_HELPER_WRITER")
.unwrap()
.parse::<u64>()
.unwrap();
let store = CheckpointStore::new(PathBuf::from(root));
for index in 0..20 {
store
.append_record(&test_record(
"shared-session",
&store.root,
writer * 1_000 + index,
))
.unwrap();
}
}
#[test]
fn cross_process_lock_preserves_all_serialized_appends() {
let temp = TempDir::new().unwrap();
let root = temp.path().join("checkpoints");
let executable = std::env::current_exe().unwrap();
let mut children = Vec::new();
for writer in 0..4 {
children.push(
std::process::Command::new(&executable)
.arg("--exact")
.arg("checkpoints::store::tests::ledger_process_append_helper")
.env("MAGI_CHECKPOINT_APPEND_HELPER_ROOT", &root)
.env("MAGI_CHECKPOINT_APPEND_HELPER_WRITER", writer.to_string())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.spawn()
.unwrap(),
);
}
for child in children {
let output = child.wait_with_output().unwrap();
assert!(
output.status.success(),
"child failed: stdout={} stderr={}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
let read = CheckpointStore::new(root).read_records("shared-session");
assert!(read.diagnostics.is_empty(), "{:?}", read.diagnostics);
assert_eq!(read.records.len(), 80);
let turns = read
.records
.into_iter()
.map(|record| record.event.user_turn)
.collect::<std::collections::BTreeSet<_>>();
assert_eq!(turns.len(), 80);
}
}