use std::fmt;
use std::fs::{File, Metadata, Permissions};
use std::io::{Read, Write};
use std::os::unix::fs::{MetadataExt, PermissionsExt};
use asupersync::Cx;
use fastmcp_core::crypto::{draw_security_identifier, sha256_bounded};
use rustix::fs::{
AtFlags, FlockOperation, Mode, OFlags, flock, openat, renameat, statat, unlinkat,
};
use rustix::io::Errno;
pub mod slot;
pub const MAX_ATOMIC_FILE_BYTES: usize = 1024 * 1024;
const MAX_LEAF_BYTES: usize = 96;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AtomicFileVersion([u8; 32]);
pub struct AtomicFileSnapshot {
version: AtomicFileVersion,
bytes: Vec<u8>,
}
impl AtomicFileSnapshot {
pub fn version(&self) -> AtomicFileVersion {
self.version
}
pub fn bytes(&self) -> &[u8] {
&self.bytes
}
pub fn into_bytes(self) -> Vec<u8> {
self.bytes
}
}
impl fmt::Debug for AtomicFileSnapshot {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("AtomicFileSnapshot")
.field("version", &self.version)
.field("byte_len", &self.bytes.len())
.finish()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum AtomicFileError {
InvalidName,
InvalidLimit,
UnsafeDirectory,
UnsafeFile,
LockReplaced,
Busy,
TooLarge,
Conflict,
RandomUnavailable,
EntropyUnavailable,
Cancelled,
TimedOut,
Io,
RecoveryRequired,
CommitUncertain {
attempted: AtomicFileVersion,
},
}
impl fmt::Display for AtomicFileError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Self::InvalidName => "atomic file requires one bounded ASCII leaf name",
Self::InvalidLimit => "atomic file byte limit is invalid",
Self::UnsafeDirectory => "atomic file directory is not owner-private",
Self::UnsafeFile => "atomic file is not an owner-private single-link regular file",
Self::LockReplaced => "atomic file lock identity changed",
Self::Busy => "atomic file is owned by another writer",
Self::TooLarge => "atomic file exceeds its configured byte limit",
Self::Conflict => "atomic file changed since the supplied snapshot",
Self::RandomUnavailable => "atomic file temporary-name randomness unavailable",
Self::EntropyUnavailable => "atomic file caller lacks entropy authority",
Self::Cancelled => "atomic file operation cancelled before commit",
Self::TimedOut => "atomic file operation deadline exceeded before commit",
Self::Io => "atomic file filesystem operation failed",
Self::RecoveryRequired => "atomic file requires explicit commit reconciliation",
Self::CommitUncertain { .. } => {
"atomic file rename completed but durability is uncertain"
}
})
}
}
impl std::error::Error for AtomicFileError {}
pub struct SecureAtomicFile {
directory: File,
lock: File,
leaf: String,
lock_leaf: String,
maximum_bytes: usize,
recovery_required: bool,
directory_sync: fn(&File) -> std::io::Result<()>,
}
impl SecureAtomicFile {
pub fn open(
cx: &Cx,
directory: File,
leaf: &str,
maximum_bytes: usize,
) -> Result<Self, AtomicFileError> {
checkpoint(cx)?;
validate_name(leaf)?;
if maximum_bytes == 0 || maximum_bytes > MAX_ATOMIC_FILE_BYTES {
return Err(AtomicFileError::InvalidLimit);
}
validate_directory(&directory)?;
let lock_leaf = format!(".{leaf}.lock");
let (lock, created) = match openat(
&directory,
lock_leaf.as_str(),
private_open_flags() | OFlags::RDWR | OFlags::CREATE | OFlags::EXCL,
Mode::RUSR | Mode::WUSR,
) {
Ok(fd) => (File::from(fd), true),
Err(Errno::EXIST) => (
File::from(
openat(
&directory,
lock_leaf.as_str(),
private_open_flags() | OFlags::RDWR,
Mode::empty(),
)
.map_err(|_| AtomicFileError::UnsafeFile)?,
),
false,
),
Err(_) => return Err(AtomicFileError::Io),
};
if created {
lock.set_permissions(Permissions::from_mode(0o600))
.map_err(|_| AtomicFileError::Io)?;
}
validate_regular(&lock.metadata().map_err(|_| AtomicFileError::Io)?)?;
checkpoint(cx)?;
match flock(&lock, FlockOperation::NonBlockingLockExclusive) {
Ok(()) => {}
Err(Errno::WOULDBLOCK) => return Err(AtomicFileError::Busy),
Err(_) => return Err(AtomicFileError::Io),
}
let store = Self {
directory,
lock,
leaf: leaf.to_owned(),
lock_leaf,
maximum_bytes,
recovery_required: false,
directory_sync: File::sync_all,
};
store.validate_authority()?;
checkpoint(cx)?;
store.lock.sync_all().map_err(|_| AtomicFileError::Io)?;
checkpoint(cx)?;
store.sync_directory().map_err(|_| AtomicFileError::Io)?;
store.read_current(cx)?;
Ok(store)
}
pub fn maximum_bytes(&self) -> usize {
self.maximum_bytes
}
pub fn load(&self, cx: &Cx) -> Result<Option<AtomicFileSnapshot>, AtomicFileError> {
self.admit(cx)?;
self.read_current(cx)
}
pub fn replace(
&mut self,
cx: &Cx,
expected: Option<AtomicFileVersion>,
bytes: &[u8],
) -> Result<AtomicFileVersion, AtomicFileError> {
self.admit(cx)?;
if bytes.len() > self.maximum_bytes {
return Err(AtomicFileError::TooLarge);
}
self.check_expected(cx, expected)?;
let version = version(bytes)?;
if !cx.capabilities().entropy {
return Err(AtomicFileError::EntropyUnavailable);
}
let mut temporary = self.create_temporary(cx)?;
write_bounded(cx, &mut temporary.file, bytes)?;
checkpoint(cx)?;
temporary.file.sync_all().map_err(|_| AtomicFileError::Io)?;
self.validate_authority()?;
self.check_expected(cx, expected)?;
checkpoint(cx)?;
renameat(
&self.directory,
temporary.leaf.as_str(),
&self.directory,
self.leaf.as_str(),
)
.map_err(|_| AtomicFileError::Io)?;
temporary.renamed = true;
drop(temporary);
if self.sync_directory().is_err() {
self.recovery_required = true;
return Err(AtomicFileError::CommitUncertain { attempted: version });
}
Ok(version)
}
fn sync_directory(&self) -> std::io::Result<()> {
(self.directory_sync)(&self.directory)
}
pub fn reconcile(&mut self, cx: &Cx) -> Result<Option<AtomicFileSnapshot>, AtomicFileError> {
checkpoint(cx)?;
self.validate_authority()?;
let current = self.read_current(cx)?;
checkpoint(cx)?;
self.sync_directory()
.map_err(|_| AtomicFileError::RecoveryRequired)?;
self.recovery_required = false;
Ok(current)
}
fn admit(&self, cx: &Cx) -> Result<(), AtomicFileError> {
checkpoint(cx)?;
if self.recovery_required {
return Err(AtomicFileError::RecoveryRequired);
}
self.validate_authority()
}
fn validate_authority(&self) -> Result<(), AtomicFileError> {
validate_directory(&self.directory)?;
let held = self.lock.metadata().map_err(|_| AtomicFileError::Io)?;
validate_regular(&held)?;
let named = statat(
&self.directory,
self.lock_leaf.as_str(),
AtFlags::SYMLINK_NOFOLLOW,
)
.map_err(|_| AtomicFileError::LockReplaced)?;
if named.st_dev != held.dev() || named.st_ino != held.ino() {
return Err(AtomicFileError::LockReplaced);
}
Ok(())
}
fn check_expected(
&self,
cx: &Cx,
expected: Option<AtomicFileVersion>,
) -> Result<(), AtomicFileError> {
let observed = self.read_current(cx)?.map(|snapshot| snapshot.version);
if observed != expected {
return Err(AtomicFileError::Conflict);
}
Ok(())
}
fn read_current(&self, cx: &Cx) -> Result<Option<AtomicFileSnapshot>, AtomicFileError> {
checkpoint(cx)?;
let mut file = match openat(
&self.directory,
self.leaf.as_str(),
private_open_flags() | OFlags::RDONLY,
Mode::empty(),
) {
Ok(fd) => File::from(fd),
Err(Errno::NOENT) => return Ok(None),
Err(Errno::LOOP) => return Err(AtomicFileError::UnsafeFile),
Err(_) => return Err(AtomicFileError::Io),
};
let metadata = file.metadata().map_err(|_| AtomicFileError::Io)?;
validate_regular(&metadata)?;
if metadata.len() > self.maximum_bytes as u64 {
return Err(AtomicFileError::TooLarge);
}
let mut bytes = Vec::with_capacity(metadata.len() as usize);
let mut buffer = [0_u8; 8192];
loop {
checkpoint(cx)?;
let remaining = self.maximum_bytes - bytes.len();
let read_limit = buffer.len().min(remaining.saturating_add(1));
let count = match file.read(&mut buffer[..read_limit]) {
Ok(0) => break,
Ok(count) => count,
Err(error) if error.kind() == std::io::ErrorKind::Interrupted => continue,
Err(_) => return Err(AtomicFileError::Io),
};
if count > remaining {
return Err(AtomicFileError::TooLarge);
}
bytes.extend_from_slice(&buffer[..count]);
}
checkpoint(cx)?;
let after = file.metadata().map_err(|_| AtomicFileError::Io)?;
validate_regular(&after)?;
if after.len() != bytes.len() as u64 || metadata.len() != after.len() {
return Err(AtomicFileError::Conflict);
}
Ok(Some(AtomicFileSnapshot {
version: version(&bytes)?,
bytes,
}))
}
fn create_temporary<'a>(&'a self, cx: &Cx) -> Result<Temporary<'a>, AtomicFileError> {
for _ in 0..4 {
checkpoint(cx)?;
let random =
draw_security_identifier().map_err(|_| AtomicFileError::RandomUnavailable)?;
let mut leaf = format!(".{}.tmp-", self.leaf);
for byte in random.as_bytes() {
use std::fmt::Write as _;
write!(&mut leaf, "{byte:02x}").map_err(|_| AtomicFileError::Io)?;
}
match openat(
&self.directory,
leaf.as_str(),
private_open_flags() | OFlags::WRONLY | OFlags::CREATE | OFlags::EXCL,
Mode::RUSR | Mode::WUSR,
) {
Ok(fd) => {
let temporary = Temporary {
directory: &self.directory,
file: File::from(fd),
leaf,
renamed: false,
};
temporary
.file
.set_permissions(Permissions::from_mode(0o600))
.map_err(|_| AtomicFileError::Io)?;
validate_regular(&temporary.file.metadata().map_err(|_| AtomicFileError::Io)?)?;
return Ok(temporary);
}
Err(Errno::EXIST) => continue,
Err(_) => return Err(AtomicFileError::Io),
}
}
Err(AtomicFileError::RandomUnavailable)
}
}
struct Temporary<'a> {
directory: &'a File,
file: File,
leaf: String,
renamed: bool,
}
impl Drop for SecureAtomicFile {
fn drop(&mut self) {
let _ = flock(&self.lock, FlockOperation::Unlock);
}
}
impl Drop for Temporary<'_> {
fn drop(&mut self) {
if !self.renamed {
if let (Ok(held), Ok(named)) = (
self.file.metadata(),
statat(
self.directory,
self.leaf.as_str(),
AtFlags::SYMLINK_NOFOLLOW,
),
) {
if named.st_dev == held.dev() && named.st_ino == held.ino() {
let _ = unlinkat(self.directory, self.leaf.as_str(), AtFlags::empty());
}
}
}
}
}
fn private_open_flags() -> OFlags {
OFlags::CLOEXEC | OFlags::NOFOLLOW | OFlags::NONBLOCK
}
fn validate_name(name: &str) -> Result<(), AtomicFileError> {
if name.is_empty()
|| name.len() > MAX_LEAF_BYTES
|| !name
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_'))
{
return Err(AtomicFileError::InvalidName);
}
Ok(())
}
fn validate_directory(directory: &File) -> Result<(), AtomicFileError> {
let metadata = directory.metadata().map_err(|_| AtomicFileError::Io)?;
if !metadata.is_dir()
|| metadata.uid() != rustix::process::geteuid().as_raw()
|| metadata.mode() & 0o7777 != 0o700
{
return Err(AtomicFileError::UnsafeDirectory);
}
Ok(())
}
fn validate_regular(metadata: &Metadata) -> Result<(), AtomicFileError> {
if !metadata.is_file()
|| metadata.uid() != rustix::process::geteuid().as_raw()
|| metadata.mode() & 0o7777 != 0o600
|| metadata.nlink() != 1
{
return Err(AtomicFileError::UnsafeFile);
}
Ok(())
}
fn version(bytes: &[u8]) -> Result<AtomicFileVersion, AtomicFileError> {
sha256_bounded(bytes, MAX_ATOMIC_FILE_BYTES)
.map(|digest| AtomicFileVersion(digest.into_bytes()))
.map_err(|_| AtomicFileError::TooLarge)
}
fn write_bounded(cx: &Cx, file: &mut File, mut bytes: &[u8]) -> Result<(), AtomicFileError> {
while !bytes.is_empty() {
checkpoint(cx)?;
match file.write(bytes) {
Ok(0) => return Err(AtomicFileError::Io),
Ok(count) => bytes = &bytes[count..],
Err(error) if error.kind() == std::io::ErrorKind::Interrupted => {}
Err(_) => return Err(AtomicFileError::Io),
}
}
Ok(())
}
fn checkpoint(cx: &Cx) -> Result<(), AtomicFileError> {
cx.checkpoint().map_err(|error| {
use asupersync::{CancelKind, error::ErrorKind};
match cx.cancel_reason().map(|reason| reason.kind) {
Some(CancelKind::Deadline | CancelKind::Timeout) => AtomicFileError::TimedOut,
Some(_) => AtomicFileError::Cancelled,
None => match error.kind() {
ErrorKind::DeadlineExceeded | ErrorKind::CancelTimeout => AtomicFileError::TimedOut,
_ => AtomicFileError::Cancelled,
},
}
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::os::unix::fs::DirBuilderExt;
use std::sync::atomic::{AtomicU64, Ordering};
static NEXT_DIRECTORY: AtomicU64 = AtomicU64::new(0);
pub(super) struct PrivateDirectory(std::path::PathBuf);
impl PrivateDirectory {
pub(super) fn new() -> Self {
let id = NEXT_DIRECTORY.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"fastmcp-commit-uncertain-{}-{id}",
std::process::id()
));
std::fs::DirBuilder::new()
.mode(0o700)
.create(&path)
.unwrap();
Self(path)
}
pub(super) fn open(&self, cx: &Cx) -> SecureAtomicFile {
SecureAtomicFile::open(cx, File::open(&self.0).unwrap(), "credential", 4096).unwrap()
}
}
impl Drop for PrivateDirectory {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
#[test]
fn a_dropped_handle_releases_its_lock_despite_an_inherited_duplicate() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let store = directory.open(&cx);
let inherited = store.lock.try_clone().unwrap();
drop(store);
let reopened =
SecureAtomicFile::open(&cx, File::open(&directory.0).unwrap(), "credential", 4096);
assert!(
reopened.is_ok(),
"a dropped handle's slot must reopen while a duplicate lives"
);
drop(inherited);
}
#[test]
fn a_live_handle_still_refuses_a_second_opener_beside_an_inherited_duplicate() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let store = directory.open(&cx);
let inherited = store.lock.try_clone().unwrap();
let second =
SecureAtomicFile::open(&cx, File::open(&directory.0).unwrap(), "credential", 4096);
assert!(matches!(second, Err(AtomicFileError::Busy)));
drop(store);
drop(inherited);
}
pub(super) fn failing_directory_sync(_: &File) -> std::io::Result<()> {
Err(std::io::Error::other("injected directory fsync failure"))
}
#[test]
fn a_failed_post_rename_directory_sync_is_commit_uncertain() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let mut store = directory.open(&cx);
let previous = store.replace(&cx, None, b"previous").unwrap();
store.directory_sync = failing_directory_sync;
let attempted = version(b"attempted").unwrap();
assert_eq!(
store.replace(&cx, Some(previous), b"attempted"),
Err(AtomicFileError::CommitUncertain { attempted }),
);
assert!(store.recovery_required);
assert!(matches!(
store.load(&cx),
Err(AtomicFileError::RecoveryRequired)
));
assert_eq!(
store.replace(&cx, Some(attempted), b"again"),
Err(AtomicFileError::RecoveryRequired),
);
store.directory_sync = File::sync_all;
let reconciled = store.reconcile(&cx).unwrap().expect("the rename happened");
assert_eq!(reconciled.version(), attempted);
assert!(!store.recovery_required);
assert_eq!(store.load(&cx).unwrap().unwrap().bytes(), b"attempted");
}
#[test]
fn a_successful_post_rename_directory_sync_commits() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let mut store = directory.open(&cx);
let previous = store.replace(&cx, None, b"previous").unwrap();
let attempted = version(b"attempted").unwrap();
assert_eq!(
store.replace(&cx, Some(previous), b"attempted"),
Ok(attempted)
);
assert!(!store.recovery_required);
assert_eq!(store.load(&cx).unwrap().unwrap().version(), attempted);
}
fn process_sharing_the_lock_description(store: &SecureAtomicFile) -> std::process::Child {
std::process::Command::new("sleep")
.arg("60")
.stdin(std::process::Stdio::from(store.lock.try_clone().unwrap()))
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.unwrap()
}
fn open_contender(
cx: &Cx,
directory: &PrivateDirectory,
) -> Result<SecureAtomicFile, AtomicFileError> {
SecureAtomicFile::open(cx, File::open(&directory.0).unwrap(), "credential", 4096)
}
#[test]
fn a_dropped_handle_is_not_kept_busy_by_another_process_sharing_its_lock() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let mut store = directory.open(&cx);
let committed = store.replace(&cx, None, b"committed").unwrap();
let mut sharer = process_sharing_the_lock_description(&store);
drop(store);
let reopened = open_contender(&cx, &directory);
let _ = sharer.kill();
let _ = sharer.wait();
let reopened = reopened.expect("dropping the only handle must release its lock");
assert_eq!(reopened.load(&cx).unwrap().unwrap().version(), committed);
}
#[test]
fn a_live_handle_keeps_a_contender_busy_while_another_process_shares_its_lock() {
let cx = Cx::for_testing();
let directory = PrivateDirectory::new();
let mut store = directory.open(&cx);
let committed = store.replace(&cx, None, b"committed").unwrap();
let mut sharer = process_sharing_the_lock_description(&store);
let contender = open_contender(&cx, &directory);
let _ = sharer.kill();
let _ = sharer.wait();
assert!(matches!(contender, Err(AtomicFileError::Busy)));
assert!(!store.recovery_required);
assert_eq!(store.load(&cx).unwrap().unwrap().version(), committed);
let next = version(b"next").unwrap();
assert_eq!(store.replace(&cx, Some(committed), b"next"), Ok(next));
}
}