Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
227 changes: 227 additions & 0 deletions crates/ghost-common/src/atomic_file.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,227 @@
//! One durable, atomic, race-safe file replacement.
//!
//! # Why this exists in one place
//!
//! Several places persist a small file by rewriting it whole, and each had
//! grown its own copy of write-temp / fsync / rename. Every copy carried the
//! same defect: the staging file was a fixed `.tmp` sibling of the target.
//!
//! A fixed staging path is shared by every concurrent writer of that file. Two
//! writers racing on it truncate each other's contents, and once the winner has
//! renamed it away the loser's `rename` fails with `ENOENT`. Staging through a
//! path private to each write removes both races, and doing it once means there
//! is no next copy to get wrong.
//!
//! Some copies were also missing durability rather than just racing: a
//! `fs::write` followed by `rename`, with neither the file nor the directory
//! fsynced, reports success for a write a power loss can still take away.
//!
//! # What this does not do
//!
//! It makes a single write atomic. It does **not** make read-modify-write
//! atomic: a caller that loads a file, edits in memory and writes the whole
//! thing back still needs to hold a lock across all three steps, or a
//! concurrent writer's changes are lost.

use std::fs;
use std::io::Write;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};

/// Distinguishes the staging files of two writes racing on one target path.
/// Paired with the pid so separate processes cannot collide either.
static TMP_SEQ: AtomicU64 = AtomicU64::new(0);

/// A staging path private to one write, beside `target`.
///
/// For writers that stream into the staging file rather than handing over a
/// finished body — a multi-hundred-megabyte proving key, say, which must not be
/// buffered in memory just to reuse [`write_atomic`]. Pair it with
/// [`sync_parent_dir`] after the rename.
///
/// Never a fixed `.tmp` sibling: that path is shared by every concurrent
/// writer of the same target, so two writers truncate each other and the
/// loser's rename fails with `ENOENT`.
pub fn staging_path(target: &Path) -> std::path::PathBuf {
let seq = TMP_SEQ.fetch_add(1, Ordering::Relaxed);
target.with_extension(format!("tmp.{}.{seq}", std::process::id()))
}

/// fsync the directory holding `path`, so a rename into it is persisted.
///
/// Renaming is a metadata change. Without this the file's own fsync survives a
/// power loss and the rename does not, which loses the write while having
/// reported success.
///
/// Best effort: a filesystem that refuses to open a directory has still given
/// us a durable staging file and an atomic rename, which is the bulk of the
/// guarantee.
pub fn sync_parent_dir(path: &Path) {
if let Some(dir) = path.parent() {
if let Ok(d) = fs::File::open(dir) {
let _ = d.sync_all();
}
}
}

/// Replace `path` with `body`, atomically and durably.
///
/// Creates the parent directory if absent. On unix, `mode` sets the staging
/// file's permissions before the rename, so the file is never briefly readable
/// at a wider mode than intended — pass `Some(0o600)` for anything secret.
///
/// The sequence is write-temp, fsync, rename, fsync-dir: what survives power
/// loss on the filesystems this runs on. Skipping the directory fsync leaves
/// the rename itself unpersisted, which loses the write while reporting
/// success.
pub fn write_atomic(path: &Path, body: &[u8], mode: Option<u32>) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent)?;
}
}

// Staging path private to this write. See the module docs for what a
// shared one costs.
let tmp = staging_path(path);

let staged = (|| -> std::io::Result<()> {
let mut f = fs::File::create(&tmp)?;
f.write_all(body)?;

#[cfg(unix)]
if let Some(mode) = mode {
use std::os::unix::fs::PermissionsExt;
f.set_permissions(fs::Permissions::from_mode(mode))?;
}
#[cfg(not(unix))]
let _ = mode;

// Contents before the rename, or the rename can land pointing at an
// empty file.
f.sync_all()?;
drop(f);
fs::rename(&tmp, path)
})();

if staged.is_err() {
// Don't leave a staging file behind for a write that failed.
let _ = fs::remove_file(&tmp);
}
staged?;

sync_parent_dir(path);
Ok(())
}

#[cfg(test)]
mod tests {
use super::*;

fn scratch(tag: &str) -> std::path::PathBuf {
let d = std::env::temp_dir().join(format!(
"ghost-atomic-{tag}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&d).expect("scratch dir");
d
}

#[test]
fn it_creates_missing_parent_directories() {
let dir = scratch("parents");
let path = dir.join("a").join("b").join("f.json");
write_atomic(&path, b"hello", None).expect("write");
assert_eq!(fs::read(&path).unwrap(), b"hello");
let _ = fs::remove_dir_all(&dir);
}

#[test]
fn it_replaces_existing_content_wholesale() {
let dir = scratch("replace");
let path = dir.join("f.json");
write_atomic(&path, b"first-and-longer", None).expect("write 1");
write_atomic(&path, b"second", None).expect("write 2");
assert_eq!(fs::read(&path).unwrap(), b"second");
let _ = fs::remove_dir_all(&dir);
}

/// The regression this module exists for.
///
/// With a fixed `.tmp` staging path, concurrent writers truncate each
/// other's staging file and the loser's rename fails with `ENOENT`.
#[test]
fn concurrent_writers_do_not_destroy_each_others_staging_file() {
let dir = scratch("race");
let path = dir.join("contended.json");

let threads: Vec<_> = (0..16u8)
.map(|i| {
let path = path.clone();
std::thread::spawn(move || {
let body = vec![b'a' + i; 64];
write_atomic(&path, &body, Some(0o600))
})
})
.collect();
for (i, t) in threads.into_iter().enumerate() {
t.join()
.expect("writer thread")
.unwrap_or_else(|e| panic!("writer {i} failed: {e}"));
}

// Every write is all-or-nothing, so whoever landed last left exactly
// its own body — never a mixture, never an empty file.
let got = fs::read(&path).expect("target exists");
assert_eq!(got.len(), 64, "target is a whole body, not a partial write");
assert!(
got.iter().all(|b| *b == got[0]),
"target interleaves two writers' bodies"
);

let strays: Vec<_> = fs::read_dir(&dir)
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.file_name().to_string_lossy().into_owned())
.filter(|n| n.contains(".tmp."))
.collect();
assert!(strays.is_empty(), "staging files left behind: {strays:?}");

let _ = fs::remove_dir_all(&dir);
}

/// Two writes to one target must never stage through the same path.
///
/// This is the whole defect, isolated: a fixed `.tmp` sibling is shared by
/// every concurrent writer, so one truncates the other's staging file and
/// the loser's rename finds nothing there. Streaming writers (a proving key
/// too large to hold in memory) use this directly rather than
/// `write_atomic`, so it needs its own guarantee.
#[test]
fn staging_paths_are_unique_per_write() {
let target = std::path::Path::new("/tmp/some-target.bin");
let a = staging_path(target);
let b = staging_path(target);
assert_ne!(a, b, "two writes staged through the same path");
assert_ne!(a, target.to_path_buf());
// Still beside the target, so the rename stays within one filesystem —
// a staging file elsewhere would make `rename` cross devices and fail.
assert_eq!(a.parent(), target.parent());
}

#[cfg(unix)]
#[test]
fn a_secret_file_is_never_wider_than_requested() {
use std::os::unix::fs::PermissionsExt;
let dir = scratch("mode");
let path = dir.join("secret.json");
write_atomic(&path, b"secret", Some(0o600)).expect("write");
let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
assert_eq!(mode, 0o600, "secret file landed at {mode:o}");
let _ = fs::remove_dir_all(&dir);
}
}
1 change: 1 addition & 0 deletions crates/ghost-common/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@

#![deny(unreachable_pub)]

pub mod atomic_file;
pub mod batch_two_phase;
pub mod circuit_breaker;
pub mod clock;
Expand Down
35 changes: 8 additions & 27 deletions crates/ghost-consensus/src/mesh.rs
Original file line number Diff line number Diff line change
Expand Up @@ -711,28 +711,9 @@ impl PeerHighWaterMarks {
path: &std::path::Path,
snapshot: &HashMap<String, (u64, u64)>,
) -> std::io::Result<()> {
use std::io::Write;

let encoded = serde_json::to_vec(snapshot)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let tmp = path.with_extension("tmp");
{
let mut f = std::fs::File::create(&tmp)?;
f.write_all(&encoded)?;
f.sync_all()?;
}
std::fs::rename(&tmp, path)?;
if let Some(parent) = path.parent() {
// Best effort: a filesystem that refuses to open a directory still gave us a durable
// temp file and an atomic rename, which is the bulk of the guarantee.
if let Ok(dir) = std::fs::File::open(parent) {
let _ = dir.sync_all();
}
}
Ok(())
ghost_common::atomic_file::write_atomic(path, &encoded, None)
}

/// Move the mark to `sequence` unconditionally, including downwards. Only for a verified
Expand Down Expand Up @@ -1946,14 +1927,14 @@ impl MeshNetwork {
}
}

/// Atomically persist the sequence ceiling (write-temp-then-rename).
/// Atomically and durably persist the sequence ceiling.
///
/// This used to stage with a bare `fs::write` and fsync neither the file
/// nor its directory, so a power loss could take back a ceiling already
/// reported as persisted — and the ceiling exists precisely so a restart
/// cannot reuse a sequence number.
fn write_sequence_ceiling(path: &std::path::Path, ceiling: u64) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let tmp = path.with_extension("tmp");
std::fs::write(&tmp, ceiling.to_string())?;
std::fs::rename(&tmp, path)
ghost_common::atomic_file::write_atomic(path, ceiling.to_string().as_bytes(), None)
}

/// Reserve a fresh lease block once the counter nears the persisted ceiling.
Expand Down
12 changes: 9 additions & 3 deletions crates/ghost-mpc/src/params.rs
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ impl ParameterFiles {
/// S-5 SECURITY: Uses atomic write (temp file + rename) to prevent corruption
/// if the process crashes mid-write. Also uses fsync for durability.
pub fn save_parameters(path: &Path, params: &Parameters<Bls12>) -> MpcResult<()> {
let temp_path = path.with_extension("tmp");
let temp_path = ghost_common::atomic_file::staging_path(path);

// Write to temp file first
let result = (|| -> MpcResult<()> {
Expand All @@ -220,6 +220,10 @@ pub fn save_parameters(path: &Path, params: &Parameters<Bls12>) -> MpcResult<()>

// Atomic rename to target path
fs::rename(&temp_path, path)?;
// The rename is a metadata change of the directory, and needs its own
// sync: without it the file's fsync above survives a power loss and the
// rename does not, losing a write already reported as saved.
ghost_common::atomic_file::sync_parent_dir(path);

let file_size = fs::metadata(path)?.len();
info!(
Expand Down Expand Up @@ -263,7 +267,7 @@ pub fn read_parameters_from_bytes(bytes: &[u8]) -> MpcResult<Parameters<Bls12>>
///
/// S-5 SECURITY: Uses atomic write (temp file + rename) to prevent corruption.
pub fn save_verifying_key(path: &Path, vk: &VerifyingKey<Bls12>) -> MpcResult<()> {
let temp_path = path.with_extension("tmp");
let temp_path = ghost_common::atomic_file::staging_path(path);

let result = (|| -> MpcResult<()> {
let file = File::create(&temp_path)?;
Expand All @@ -283,6 +287,7 @@ pub fn save_verifying_key(path: &Path, vk: &VerifyingKey<Bls12>) -> MpcResult<()
}

fs::rename(&temp_path, path)?;
ghost_common::atomic_file::sync_parent_dir(path);

info!(path = %path.display(), "Saved verifying key (atomic write)");

Expand Down Expand Up @@ -366,7 +371,7 @@ pub fn update_current_params(files: &ParameterFiles, version: u32) -> MpcResult<

/// Atomically copy a file by writing to a temp file and renaming.
fn atomic_copy(src: &Path, dst: &Path) -> MpcResult<()> {
let temp_path = dst.with_extension("tmp");
let temp_path = ghost_common::atomic_file::staging_path(dst);

let result = (|| -> MpcResult<()> {
let mut reader = BufReader::new(File::open(src)?);
Expand All @@ -385,6 +390,7 @@ fn atomic_copy(src: &Path, dst: &Path) -> MpcResult<()> {
}

fs::rename(&temp_path, dst)?;
ghost_common::atomic_file::sync_parent_dir(dst);
Ok(())
}

Expand Down
Loading