perf(ecstore): Linux O_DIRECT shard-write path off the commit critical section (#4411)

`create_file` writes erasure shard / multipart part data through the page
cache, so the commit-point `sync_dir_files` fdatasync in `rename_data` has to
flush the whole shard's dirty pages (profiling: create_file avg 67.52ms vs
1.51ms nosync; BitrotWriter::write p95 40.86ms), and concurrent PUTs on one
device stall each other's writeback inside the rename critical section.

Add an opt-in Linux O_DIRECT streaming writer that reuses the direct-io read
infrastructure (statx DIOALIGN probe, aligned bounce buffer, supported latch,
one-time fallback log). Incoming bytes are staged into an aligned buffer;
full aligned batches and the shutdown tail are flushed on the blocking pool
(same offloading posture as the buffered writer it replaces). The
sub-alignment tail is written buffered after clearing O_DIRECT via fcntl
(MinIO's recipe). Shard bytes reach the device during the write phase, so the
unchanged commit-point `sync_dir_files` fdatasync degrades to a cheap
metadata/device FLUSH instead of flushing ~2MiB of dirty pages.

Gated by `RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE` (default false), Linux-only,
and durability-neutral: `sync_dir_files` still fdatasyncs every shard at the
commit point, so the strict/relaxed/none durability tiers keep their exact
guarantees. Filesystems that reject O_DIRECT (tmpfs, overlayfs, 9p) latch the
path off and fall back to buffered writes; an O_DIRECT open/write EINVAL is
never surfaced as `InvalidInput` (which `to_file_error` maps to `FileNotFound`
and would masquerade as a missing shard, triggering a spurious EC rebuild).

Unit tests cover the env gate, the alignment/staging-capacity and tail-split
math, an end-to-end `create_file` round trip across sizes crossing the
alignment and multi-batch boundaries (buffered fallback on macOS / O_DIRECT on
a block device), and the writer state machine over a plain file on Linux. Real
O_DIRECT throughput on Linux NVMe (ext4/xfs) is deferred to the azure 4-node
bench per the issue's GO/NO-GO gate; the development host is macOS and cannot
exercise O_DIRECT.

Refs: rustfs/backlog#927, rustfs/backlog#936

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-08 16:39:43 +08:00
committed by GitHub
parent 021c955c21
commit 92c8c6db75
+562
View File
@@ -227,6 +227,20 @@ const DEFAULT_RUSTFS_OBJECT_DIRECT_IO_READ_ENABLE: bool = false;
/// Default: 4MB.
const ENV_RUSTFS_OBJECT_DIRECT_IO_READ_THRESHOLD: &str = "RUSTFS_OBJECT_DIRECT_IO_READ_THRESHOLD";
const DEFAULT_RUSTFS_OBJECT_DIRECT_IO_READ_THRESHOLD: usize = 4 * 1024 * 1024;
/// Enable O_DIRECT for erasure shard / multipart part data writes (Linux only).
/// When enabled, `create_file` streams shard bytes straight to the device with
/// O_DIRECT, so the commit-point `sync_dir_files` fdatasync no longer flushes
/// ~2 MiB of dirty pages inside the `rename_data` critical section (it degrades
/// to a cheap metadata/device FLUSH). Aligned whole blocks are written direct;
/// the trailing sub-alignment remainder falls back to a buffered write after
/// clearing O_DIRECT (MinIO's recipe). Durability is unchanged: the file is
/// still fdatasynced at the commit point by the unchanged `sync_dir_files`.
/// EINVAL/EOPNOTSUPP (tmpfs, overlayfs, 9p, ...) latch the path off and fall
/// back to buffered writes for the whole disk. Non-Linux always falls back.
/// Default: false (buffered writes via the page cache, as before).
const ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE: &str = "RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE";
const DEFAULT_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE: bool = false;
const ENV_RUSTFS_OBJECT_MMAP_POPULATE_ENABLE: &str = "RUSTFS_OBJECT_MMAP_POPULATE_ENABLE";
const DEFAULT_RUSTFS_OBJECT_MMAP_POPULATE_ENABLE: bool = false;
const ENV_RUSTFS_OBJECT_MMAP_READ_METHOD: &str = "RUSTFS_OBJECT_MMAP_READ_METHOD";
@@ -257,6 +271,11 @@ fn is_direct_io_read_enabled() -> bool {
rustfs_utils::get_env_bool(ENV_RUSTFS_OBJECT_DIRECT_IO_READ_ENABLE, DEFAULT_RUSTFS_OBJECT_DIRECT_IO_READ_ENABLE)
}
/// Check if O_DIRECT shard/part data writes are enabled.
fn is_direct_io_write_enabled() -> bool {
rustfs_utils::get_env_bool(ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE, DEFAULT_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE)
}
const EVENT_DISK_LOCAL_DURABILITY_MODE: &str = "disk_local_durability_mode";
/// Process-wide durability tier for commit-point fsync work on the local disk.
@@ -711,6 +730,376 @@ fn pread_direct_aligned(file_path: &Path, offset: u64, length: usize, state: &Di
Ok(Bytes::copy_from_slice(&buf.as_slice()[logical_start..logical_end]))
}
// `AlignedBuf` uniquely owns a single heap allocation reached only through
// `&self`/`&mut self`; there is no interior mutability and no aliasing, so it
// is sound to move it across threads (into a `spawn_blocking` flush closure)
// and to share `&AlignedBuf` between threads.
#[cfg(target_os = "linux")]
#[allow(unsafe_code)]
// SAFETY: exclusive heap ownership, no aliasing (see the note above).
unsafe impl Send for AlignedBuf {}
#[cfg(target_os = "linux")]
#[allow(unsafe_code)]
// SAFETY: `&AlignedBuf` only exposes read-only access to an immutable buffer.
unsafe impl Sync for AlignedBuf {}
/// Runtime state for the true O_DIRECT write path (Linux only), mirroring
/// [`DirectIoReadState`].
///
/// `supported` starts true and latches false on the first EINVAL/EOPNOTSUPP
/// from an O_DIRECT open (tmpfs, overlayfs, and 9p commonly reject the flag);
/// `create_file` then permanently opens shard files buffered for this disk.
/// `align` caches the DIO alignment probed from the backing filesystem. As on
/// the read path, an O_DIRECT open error must never surface to callers: EINVAL
/// maps to `FileNotFound` in `to_file_error`, which would masquerade as a
/// missing shard and trigger spurious EC rebuilds.
#[cfg(target_os = "linux")]
#[derive(Debug)]
struct DirectIoWriteState {
supported: AtomicBool,
align: OnceLock<usize>,
fallback_logged: AtomicBool,
}
#[cfg(target_os = "linux")]
impl DirectIoWriteState {
fn new() -> Self {
Self {
supported: AtomicBool::new(true),
align: OnceLock::new(),
fallback_logged: AtomicBool::new(false),
}
}
}
/// Target staging size for O_DIRECT writes, rounded up to the DIO alignment.
/// Bounds the per-writer aligned bounce buffer and batches many shard blocks
/// into one positioned write to keep the syscall count low.
const DIRECT_WRITE_STAGING_BYTES: usize = 1024 * 1024;
/// Aligned bounce-buffer capacity for a given DIO alignment: the target staging
/// size rounded up to a whole multiple of `align` so the buffer address, every
/// flushed batch length, and every write offset stay alignment-correct.
/// Platform-independent (no O_DIRECT), so it is unit-tested on any host.
fn direct_write_staging_capacity(align: usize) -> usize {
debug_assert!(align.is_power_of_two() && align >= 512);
DIRECT_WRITE_STAGING_BYTES.div_ceil(align) * align
}
/// Split `filled` staged bytes into the alignment-sized prefix written with
/// O_DIRECT and the sub-alignment tail written buffered. Platform-independent,
/// so the tail-boundary math is unit-tested on any host.
fn direct_write_tail_split(filled: usize, align: usize) -> (usize, usize) {
let aligned = filled - (filled % align);
(aligned, filled - aligned)
}
/// Positioned write-all helper: retries short writes at increasing offsets.
///
/// Under O_DIRECT the buffer address, `offset`, and length must all be aligned;
/// callers guarantee that. After O_DIRECT has been cleared (tail path) there is
/// no alignment requirement.
#[cfg(target_os = "linux")]
fn pwrite_all(file: &std::fs::File, mut buf: &[u8], mut offset: u64) -> std::io::Result<()> {
use std::os::unix::fs::FileExt;
while !buf.is_empty() {
let n = file.write_at(buf, offset)?;
if n == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::WriteZero,
"O_DIRECT positioned write wrote 0 bytes",
));
}
buf = &buf[n..];
offset += n as u64;
}
Ok(())
}
/// Never let an O_DIRECT write error reach `to_file_error` as `InvalidInput`:
/// that maps to `FileNotFound` and would masquerade as a missing shard,
/// triggering a spurious EC rebuild (backlog#897 / issue correction #2). Any
/// EINVAL/EOPNOTSUPP surfacing from a flush is remapped to a generic error so
/// the write-quorum machinery treats it as the real write failure it is.
#[cfg(target_os = "linux")]
fn sanitize_direct_write_error(err: std::io::Error) -> std::io::Error {
if is_direct_io_unsupported(&err) {
std::io::Error::other(format!("O_DIRECT shard write failed: {err}"))
} else {
err
}
}
/// Owned O_DIRECT write state moved in and out of the `spawn_blocking` flush
/// closures so the reactor is never blocked on synchronous device I/O.
#[cfg(target_os = "linux")]
struct DirectWriteInner {
file: std::fs::File,
/// Aligned bounce buffer; its capacity (`buf.len`) is a whole multiple of
/// `align`.
buf: AlignedBuf,
/// Bytes currently staged in `buf` and not yet written to the device.
filled: usize,
/// Next file offset for an O_DIRECT positioned write; always a multiple of
/// `align` because every batch flushed before the tail is a whole multiple.
write_offset: u64,
align: usize,
direct_cleared: bool,
}
#[cfg(target_os = "linux")]
impl DirectWriteInner {
/// Flush a full staging batch (`filled == buf capacity`, a multiple of
/// `align`) straight to the device with O_DIRECT.
fn flush_batch(&mut self) -> std::io::Result<()> {
if self.filled == 0 {
return Ok(());
}
debug_assert_eq!(self.filled % self.align, 0, "batch flush must be alignment-sized");
pwrite_all(&self.file, &self.buf.as_slice()[..self.filled], self.write_offset).map_err(sanitize_direct_write_error)?;
self.write_offset += self.filled as u64;
self.filled = 0;
Ok(())
}
/// Final flush at shutdown: write the aligned prefix with O_DIRECT, then the
/// sub-alignment tail buffered after clearing O_DIRECT (MinIO's recipe; the
/// tail is not separately fsynced — the commit-point `sync_dir_files`
/// fdatasync covers the whole file, issue correction #5).
fn finish(&mut self) -> std::io::Result<()> {
let (aligned, remainder) = direct_write_tail_split(self.filled, self.align);
if aligned > 0 {
pwrite_all(&self.file, &self.buf.as_slice()[..aligned], self.write_offset).map_err(sanitize_direct_write_error)?;
self.write_offset += aligned as u64;
}
if remainder > 0 {
self.clear_direct()?;
// Snapshot the slice bounds first to avoid borrowing `self.buf`
// while `self.file` is borrowed immutably below.
let start = aligned;
let end = self.filled;
pwrite_all(&self.file, &self.buf.as_slice()[start..end], self.write_offset)?;
self.write_offset += remainder as u64;
}
self.filled = 0;
Ok(())
}
/// Drop the O_DIRECT flag from the open file so the unaligned tail can be
/// written through the page cache without an alignment fault.
fn clear_direct(&mut self) -> std::io::Result<()> {
if self.direct_cleared {
return Ok(());
}
let flags = rustix::fs::fcntl_getfl(&self.file).map_err(std::io::Error::from)?;
rustix::fs::fcntl_setfl(&self.file, flags - rustix::fs::OFlags::DIRECT).map_err(std::io::Error::from)?;
self.direct_cleared = true;
Ok(())
}
}
#[cfg(target_os = "linux")]
type DirectFlushHandle = tokio::task::JoinHandle<(DirectWriteInner, std::io::Result<()>)>;
#[cfg(target_os = "linux")]
enum DirectWriteState {
Idle(Option<DirectWriteInner>),
Busy(DirectFlushHandle),
}
/// Streaming O_DIRECT writer returned by `create_file` on Linux when the path
/// is enabled and supported.
///
/// Incoming bytes are memcpy'd into an aligned bounce buffer (cheap, on the
/// reactor); each full aligned batch and the shutdown tail are flushed on the
/// blocking pool so the reactor never stalls on synchronous device I/O — the
/// same offloading posture as the buffered `tokio::fs::File` writer it
/// replaces. Durability is unchanged: no fsync happens here; the commit-point
/// `sync_dir_files` fdatasync persists the file.
#[cfg(target_os = "linux")]
struct DirectWriter {
state: DirectWriteState,
shutdown_started: bool,
shutdown_done: bool,
}
#[cfg(target_os = "linux")]
impl DirectWriter {
fn new(inner: DirectWriteInner) -> Self {
Self {
state: DirectWriteState::Idle(Some(inner)),
shutdown_started: false,
shutdown_done: false,
}
}
/// Build a writer over an already-open plain file with a caller-chosen
/// alignment and staging capacity. Exercises the streaming/tail state
/// machine deterministically on CI filesystems that reject O_DIRECT
/// (tmpfs/overlayfs), where the production `open_direct_writer` would latch
/// off; `write_at`/`fcntl` behave identically on a buffered file.
#[cfg(test)]
fn from_std_file_for_test(file: std::fs::File, align: usize, capacity: usize) -> Self {
assert_eq!(capacity % align, 0, "test capacity must be an alignment multiple");
let buf = AlignedBuf::new(capacity, align).expect("aligned buffer allocation");
Self::new(DirectWriteInner {
file,
buf,
filled: 0,
write_offset: 0,
align,
direct_cleared: false,
})
}
/// Drive an in-flight flush to completion, returning the recovered inner
/// state. Returns `Pending`/errors verbatim; on success the state is left
/// `Idle`.
fn poll_drive_busy(&mut self, cx: &mut std::task::Context<'_>) -> std::task::Poll<std::io::Result<()>> {
if let DirectWriteState::Busy(handle) = &mut self.state {
let (inner, res) = match std::task::ready!(std::pin::Pin::new(handle).poll(cx)) {
Ok(pair) => pair,
Err(join_err) => {
return std::task::Poll::Ready(Err(std::io::Error::other(format!("O_DIRECT flush task failed: {join_err}"))));
}
};
self.state = DirectWriteState::Idle(Some(inner));
if self.shutdown_started {
self.shutdown_done = true;
}
res?;
}
std::task::Poll::Ready(Ok(()))
}
}
#[cfg(target_os = "linux")]
impl AsyncWrite for DirectWriter {
fn poll_write(
self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
buf: &[u8],
) -> std::task::Poll<std::io::Result<usize>> {
let this = self.get_mut();
loop {
match &mut this.state {
DirectWriteState::Busy(_) => {
std::task::ready!(this.poll_drive_busy(cx))?;
}
DirectWriteState::Idle(inner_opt) => {
if buf.is_empty() {
return std::task::Poll::Ready(Ok(0));
}
let inner = inner_opt.as_mut().expect("idle direct writer must hold inner state");
let capacity = inner.buf.len;
let space = capacity - inner.filled;
let n = space.min(buf.len());
let start = inner.filled;
inner.buf.as_mut_slice()[start..start + n].copy_from_slice(&buf[..n]);
inner.filled += n;
if inner.filled == capacity {
let mut inner = inner_opt.take().expect("idle direct writer must hold inner state");
let handle = tokio::task::spawn_blocking(move || {
let res = inner.flush_batch();
(inner, res)
});
this.state = DirectWriteState::Busy(handle);
}
return std::task::Poll::Ready(Ok(n));
}
}
}
}
fn poll_flush(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> std::task::Poll<std::io::Result<()>> {
// Only drive an in-flight batch to completion. Sub-alignment staged
// bytes cannot be flushed mid-stream (they would misalign the next
// O_DIRECT offset); they are written by `poll_shutdown`.
self.get_mut().poll_drive_busy(cx)
}
fn poll_shutdown(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> std::task::Poll<std::io::Result<()>> {
let this = self.get_mut();
loop {
match &mut this.state {
DirectWriteState::Busy(_) => {
std::task::ready!(this.poll_drive_busy(cx))?;
}
DirectWriteState::Idle(inner_opt) => {
if this.shutdown_done {
return std::task::Poll::Ready(Ok(()));
}
let mut inner = inner_opt.take().expect("idle direct writer must hold inner state");
this.shutdown_started = true;
let handle = tokio::task::spawn_blocking(move || {
let res = inner.finish();
(inner, res)
});
this.state = DirectWriteState::Busy(handle);
}
}
}
}
}
/// Open a shard file for an O_DIRECT streaming write, probing DIO alignment on
/// the freshly created file. Returns `Ok(None)` (with the state latched off and
/// a one-time warning) when the filesystem rejects O_DIRECT, so the caller can
/// fall back to the buffered writer without ever surfacing EINVAL.
#[cfg(target_os = "linux")]
fn open_direct_writer(file_path: &Path, state: &DirectIoWriteState) -> Result<Option<DirectWriter>> {
use std::os::unix::fs::OpenOptionsExt;
let open_result = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.custom_flags(libc::O_DIRECT)
.open(file_path);
let file = match open_result {
Ok(file) => file,
Err(err) => {
if is_direct_io_unsupported(&err) {
state.supported.store(false, Ordering::Relaxed);
if !state.fallback_logged.swap(true, Ordering::Relaxed) {
warn!(
event = EVENT_DISK_LOCAL_DIRECT_IO_FALLBACK,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_DISK_LOCAL,
path = %file_path.display(),
error = ?err,
"O_DIRECT write unavailable; falling back to buffered writes"
);
}
return Ok(None);
}
// A genuine open failure (permissions, disk full, ...): map it as a
// normal file error so callers see the real cause.
return Err(to_file_error(err).into());
}
};
let align = *state.align.get_or_init(|| probe_direct_io_align(&file));
let capacity = direct_write_staging_capacity(align);
let buf = match AlignedBuf::new(capacity, align) {
Ok(buf) => buf,
Err(err) => return Err(to_file_error(err).into()),
};
Ok(Some(DirectWriter::new(DirectWriteInner {
file,
buf,
filled: 0,
write_offset: 0,
align,
direct_cleared: false,
})))
}
#[cfg(unix)]
#[allow(unsafe_code)]
fn mmap_page_size() -> Result<u64> {
@@ -1202,6 +1591,8 @@ pub(crate) struct StdBackend {
root: PathBuf,
#[cfg(target_os = "linux")]
direct_io: Arc<DirectIoReadState>,
#[cfg(target_os = "linux")]
direct_io_write: Arc<DirectIoWriteState>,
}
impl StdBackend {
@@ -1210,6 +1601,8 @@ impl StdBackend {
root,
#[cfg(target_os = "linux")]
direct_io: Arc::new(DirectIoReadState::new()),
#[cfg(target_os = "linux")]
direct_io_write: Arc::new(DirectIoWriteState::new()),
}
}
@@ -1752,6 +2145,25 @@ impl LocalIoBackend for StdBackend {
if let Some(parent) = file_path.parent() {
os::make_dir_all(parent, &volume_dir).await?;
}
// O_DIRECT streaming write (Linux, opt-in): shard bytes stream
// straight to the device so the commit-point fdatasync no longer
// flushes the whole shard's dirty pages inside the rename_data
// critical section. Latches off and falls back to the buffered
// writer on filesystems that reject O_DIRECT; never surfaces the
// EINVAL (which would masquerade as a missing shard).
#[cfg(target_os = "linux")]
if is_direct_io_write_enabled() && self.direct_io_write.supported.load(Ordering::Relaxed) {
let write_state = self.direct_io_write.clone();
let direct_path = file_path.clone();
let direct = tokio::task::spawn_blocking(move || open_direct_writer(&direct_path, &write_state))
.await
.map_err(|err| DiskError::other(format!("O_DIRECT open task failed: {err}")))??;
if let Some(writer) = direct {
return Ok(Box::new(writer));
}
}
// O_TRUNC: if a file already exists at this path, stale trailing bytes past
// the new content would otherwise survive and mismatch the metadata size.
let f = super::fs::open_file(&file_path, O_CREATE | O_WRONLY | O_TRUNC)
@@ -7994,6 +8406,156 @@ mod test {
assert_eq!(result.results, vec![CHECK_PART_FILE_CORRUPT]);
}
// ----- HP-6: O_DIRECT shard write -----
#[test]
fn direct_io_write_env_gate_defaults_off() {
temp_env::with_var_unset(ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE, || {
assert!(!is_direct_io_write_enabled(), "O_DIRECT write must default to off");
});
temp_env::with_var(ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE, Some("true"), || {
assert!(is_direct_io_write_enabled());
});
temp_env::with_var(ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE, Some("false"), || {
assert!(!is_direct_io_write_enabled());
});
}
#[test]
fn direct_write_staging_capacity_is_smallest_alignment_multiple_covering_target() {
for align in [512usize, 1024, 4096, 8192, 65536] {
let cap = direct_write_staging_capacity(align);
assert_eq!(cap % align, 0, "capacity must be an alignment multiple (align={align})");
assert!(
cap >= DIRECT_WRITE_STAGING_BYTES,
"capacity must cover the target staging size (align={align})"
);
assert!(
cap - DIRECT_WRITE_STAGING_BYTES < align,
"capacity must be the smallest covering multiple (align={align})"
);
}
}
#[test]
fn direct_write_tail_split_separates_aligned_prefix_from_remainder() {
let align = 4096usize;
assert_eq!(direct_write_tail_split(0, align), (0, 0));
assert_eq!(direct_write_tail_split(align, align), (align, 0));
assert_eq!(direct_write_tail_split(3 * align, align), (3 * align, 0));
assert_eq!(direct_write_tail_split(100, align), (0, 100));
let (prefix, remainder) = direct_write_tail_split(2 * align + 123, align);
assert_eq!((prefix, remainder), (2 * align, 123));
assert_eq!(prefix + remainder, 2 * align + 123, "split must reconstruct the input length");
assert_eq!(prefix % align, 0, "prefix must be alignment-sized");
assert!(remainder < align, "remainder must be sub-alignment");
}
/// End-to-end round trip through `create_file`'s writer with O_DIRECT writes
/// enabled. On a block-backed Linux filesystem this drives the true
/// O_DIRECT path; on macOS and on CI filesystems that reject O_DIRECT it
/// exercises the buffered fallback. Either way, enabling the gate must not
/// change the bytes read back, across sizes crossing the alignment and
/// multi-batch staging boundaries (zero-breakage contract).
#[cfg(unix)]
#[tokio::test]
async fn create_file_direct_write_round_trips_all_sizes() {
use tempfile::tempdir;
let root_dir = tempdir().expect("tempdir");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint");
let disk = LocalDisk::new(&endpoint, false).await.expect("disk");
disk.make_volume("test-volume").await.expect("make_volume");
let sizes = [
1usize,
511,
512,
4095,
4096,
4097,
8192,
12345,
DIRECT_WRITE_STAGING_BYTES - 1,
DIRECT_WRITE_STAGING_BYTES,
DIRECT_WRITE_STAGING_BYTES + 4096,
2 * DIRECT_WRITE_STAGING_BYTES + 777,
];
for (i, size) in sizes.into_iter().enumerate() {
let mut state = 0x1234_5678u64.wrapping_add(i as u64);
let content: Vec<u8> = (0..size)
.map(|_| {
state = state.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
(state >> 33) as u8
})
.collect();
let path = format!("obj-{i}.bin");
temp_env::async_with_vars([(ENV_RUSTFS_OBJECT_DIRECT_IO_WRITE_ENABLE, Some("true"))], async {
let mut writer = disk
.create_file("", "test-volume", &path, size as i64)
.await
.expect("create_file");
// Odd-sized chunks exercise partial staging fills and flushes.
let mut off = 0;
while off < content.len() {
let end = (off + 777).min(content.len());
writer.write_all(&content[off..end]).await.expect("write_all");
off = end;
}
writer.shutdown().await.expect("shutdown");
})
.await;
let got = disk.read_file_mmap_copy("test-volume", &path, 0, size).await.expect("read");
assert_eq!(got.as_ref(), content.as_slice(), "round-trip mismatch at size={size}");
}
}
/// Deterministic coverage of the O_DIRECT writer's aligned-batch + tail
/// state machine on Linux, independent of whether the CI filesystem
/// supports O_DIRECT: the writer is driven over a plain file with a small
/// alignment and staging capacity so both a full-batch flush and the tail
/// split are exercised for every size class.
#[cfg(target_os = "linux")]
#[tokio::test]
async fn direct_writer_state_machine_round_trips_over_plain_file() {
use std::io::Read;
use tempfile::tempdir;
let dir = tempdir().expect("tempdir");
let align = 512usize;
let capacity = 4 * align; // forces multi-batch flushes for larger sizes
for &size in &[0usize, 1, 511, 512, 513, 1024, 2048, 2049, 5000] {
let path = dir.path().join(format!("s{size}"));
let file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.read(true)
.open(&path)
.expect("open plain file");
let content: Vec<u8> = (0..size).map(|i| (i.wrapping_mul(7).wrapping_add(3)) as u8).collect();
let mut writer = DirectWriter::from_std_file_for_test(file, align, capacity);
let mut off = 0;
while off < content.len() {
let end = (off + 300).min(content.len());
writer.write_all(&content[off..end]).await.expect("write_all");
off = end;
}
writer.shutdown().await.expect("shutdown");
let mut got = Vec::new();
std::fs::File::open(&path)
.expect("reopen")
.read_to_end(&mut got)
.expect("read_to_end");
assert_eq!(got, content, "state-machine round-trip mismatch at size={size}");
}
}
/// Differential test for the LocalIoBackend refactor: every read shape
/// (pread via mmap_copy, pread via direct_read_copy, bounded stream, full
/// read) must return byte-identical data for the same (offset, length),