use crate::platform::traits::{DeviceReader, DeviceWriter};
use anyhow::{Context, Result};
use log::debug;
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::os::unix::io::AsRawFd;
pub(crate) const SYNC_CHUNK_BYTES: u64 = 4 * 1024 * 1024;
const SYNC_FILE_RANGE_WRITE: u32 = 2;
const SYNC_FILE_RANGE_WAIT_AFTER: u32 = 4;
fn sync_file_range_chunks(
fd: std::os::unix::io::RawFd,
offset: u64,
len: u64,
mut on_progress: Option<&mut dyn FnMut(u64, u64)>,
) -> Result<()> {
if len == 0 {
return Ok(());
}
let flags = SYNC_FILE_RANGE_WRITE | SYNC_FILE_RANGE_WAIT_AFTER;
let mut synced_in_range = 0u64;
while synced_in_range < len {
let chunk = SYNC_CHUNK_BYTES.min(len - synced_in_range);
let chunk_i64 = i64::try_from(chunk).context("Sync chunk too large")?;
let offset_i64 = i64::try_from(offset + synced_in_range).context("Sync offset too large")?;
let result = unsafe { libc::sync_file_range(fd, offset_i64, chunk_i64, flags) };
if result != 0 {
debug!(
"sync_file_range failed at offset {}: {}; falling back to fsync",
offset + synced_in_range,
std::io::Error::last_os_error()
);
unsafe {
if libc::fsync(fd) != 0 {
return Err(std::io::Error::last_os_error()).context("Failed to sync device");
}
}
if let Some(callback) = on_progress.as_mut() {
callback(len, len);
}
return Ok(());
}
synced_in_range += chunk;
if let Some(callback) = on_progress.as_mut() {
callback(synced_in_range, len);
}
}
Ok(())
}
pub struct LinuxDeviceReader {
file: File,
}
pub struct LinuxBufferedDeviceReader {
file: File,
}
impl DeviceReader for LinuxBufferedDeviceReader {
fn open(device_path: &str) -> Result<Self> {
debug!(
"Opening Linux device for buffered clone/verify read: {}",
device_path
);
let file = OpenOptions::new()
.read(true)
.open(device_path)
.context(format!(
"Failed to open device for verification read: {}",
device_path
))?;
Ok(Self { file })
}
fn device_size(&self) -> Result<u64> {
let metadata = self
.file
.metadata()
.context("Failed to get device metadata")?;
Ok(metadata.len())
}
}
impl Read for LinuxBufferedDeviceReader {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
self.file.read(buf)
}
}
impl DeviceReader for LinuxDeviceReader {
fn open(device_path: &str) -> Result<Self> {
debug!("Opening Linux device for reading: {}", device_path);
let file = OpenOptions::new()
.read(true)
.open(device_path)
.context(format!(
"Failed to open device for reading: {}",
device_path
))?;
Ok(Self { file })
}
fn device_size(&self) -> Result<u64> {
let metadata = self
.file
.metadata()
.context("Failed to get device metadata")?;
Ok(metadata.len())
}
}
impl Read for LinuxDeviceReader {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
self.file.read(buf)
}
}
pub struct LinuxDeviceWriter {
file: File,
}
impl DeviceWriter for LinuxDeviceWriter {
fn open(device_path: &str) -> Result<Self> {
debug!("Opening Linux device for writing: {}", device_path);
let file = OpenOptions::new()
.read(true)
.write(true)
.open(device_path)
.context(format!(
"Failed to open device for writing: {}",
device_path
))?;
Ok(Self { file })
}
fn write_at(&mut self, offset: u64, buf: &[u8]) -> Result<()> {
self.file
.seek(SeekFrom::Start(offset))
.context(format!("Failed to seek device to offset {offset}"))?;
self.file
.write_all(buf)
.context(format!("Failed to write {} bytes at offset {offset}", buf.len()))?;
Ok(())
}
fn supports_incremental_sync(&self) -> bool {
true
}
fn sync_written_range(&mut self, offset: u64, len: u64) -> Result<()> {
if len == 0 {
return Ok(());
}
self.file.flush().context("Failed to flush device")?;
sync_file_range_chunks(self.file.as_raw_fd(), offset, len, None)?;
debug!("Synced {len} bytes at device offset {offset}");
Ok(())
}
fn flush_and_sync(&mut self) -> Result<()> {
self.flush_and_sync_with_progress(0, None)
}
fn flush_and_sync_with_progress(
&mut self,
sync_bytes: u64,
on_progress: Option<&mut dyn FnMut(u64, u64)>,
) -> Result<()> {
self.file.flush().context("Failed to flush device")?;
if sync_bytes == 0 {
unsafe {
if libc::fsync(self.file.as_raw_fd()) != 0 {
return Err(std::io::Error::last_os_error()).context("Failed to sync device");
}
}
debug!("Device flushed and synced");
return Ok(());
}
sync_file_range_chunks(self.file.as_raw_fd(), 0, sync_bytes, on_progress)?;
debug!("Device flushed and synced ({sync_bytes} bytes)");
Ok(())
}
fn device_size(&self) -> Result<u64> {
let metadata = self
.file
.metadata()
.context("Failed to get device metadata")?;
Ok(metadata.len())
}
fn supports_inline_verify(&self) -> bool {
true
}
fn rewind_for_verify(&mut self) -> Result<()> {
self.file
.seek(SeekFrom::Start(0))
.context("Failed to rewind device for verification")?;
Ok(())
}
fn read_for_verify(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
self.file.read(buf)
}
}
impl Write for LinuxDeviceWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.file.write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
self.file.flush()
}
}