Skip to main content

mpi/
io.rs

1//! Parallel file I/O ([`File`]), corresponding to MPI's I/O chapter
2//! (`MPI_File_*`), which released rsmpi does not currently expose.
3//!
4//! Each rank opens the same path with its own OS file handle; explicit-offset
5//! reads and writes seek to a byte offset and transfer a typed buffer.
6//! Collective variants add a barrier so all ranks participate together. Offsets
7//! are measured in **bytes**.
8//!
9//! ```no_run
10//! use mpi::traits::*;
11//! use mpi::io::{File, MODE_CREATE, MODE_RDWR};
12//!
13//! let universe = mpi::initialize().unwrap();
14//! let world = universe.world();
15//! let mut f = File::open(&world, "/tmp/data.bin", MODE_CREATE | MODE_RDWR).unwrap();
16//! let rank = world.rank();
17//! // each rank writes its block at a distinct offset
18//! f.write_at_all((rank as u64) * 16, &[rank; 4]);
19//! ```
20
21use std::fs::OpenOptions;
22use std::io::{self, Read, Seek, SeekFrom, Write};
23use std::path::{Path, PathBuf};
24
25use crate::collective::CommunicatorCollectives;
26use crate::datatype::Equivalence;
27use crate::topology::{Communicator, SimpleCommunicator};
28
29/// Open the file for reading (`MPI_MODE_RDONLY`).
30pub const MODE_RDONLY: u32 = 1;
31/// Open the file for writing (`MPI_MODE_WRONLY`).
32pub const MODE_WRONLY: u32 = 2;
33/// Open the file for reading and writing (`MPI_MODE_RDWR`).
34pub const MODE_RDWR: u32 = 4;
35/// Create the file if it does not exist (`MPI_MODE_CREATE`).
36pub const MODE_CREATE: u32 = 8;
37/// Open in append mode (`MPI_MODE_APPEND`).
38pub const MODE_APPEND: u32 = 16;
39/// Fail if the file already exists (`MPI_MODE_EXCL`).
40pub const MODE_EXCL: u32 = 32;
41/// Delete the file when it is closed (`MPI_MODE_DELETE_ON_CLOSE`).
42pub const MODE_DELETE_ON_CLOSE: u32 = 64;
43
44#[inline]
45fn as_bytes<T>(v: &[T]) -> &[u8] {
46    // SAFETY: callers use POD `Equivalence` types.
47    unsafe { std::slice::from_raw_parts(v.as_ptr() as *const u8, std::mem::size_of_val(v)) }
48}
49
50#[inline]
51fn as_bytes_mut<T>(v: &mut [T]) -> &mut [u8] {
52    // SAFETY: as above.
53    unsafe { std::slice::from_raw_parts_mut(v.as_mut_ptr() as *mut u8, std::mem::size_of_val(v)) }
54}
55
56/// A file opened collectively by a communicator (`MPI_File`).
57pub struct File {
58    inner: std::fs::File,
59    comm: SimpleCommunicator,
60    path: PathBuf,
61    delete_on_close: bool,
62}
63
64impl File {
65    /// Collectively open `path` with the given mode flags (`MPI_File_open`).
66    pub fn open<C: Communicator>(comm: &C, path: impl AsRef<Path>, mode: u32) -> io::Result<File> {
67        let path = path.as_ref().to_path_buf();
68
69        // Rank 0 handles creation so CREATE/EXCL are well-defined.
70        if mode & MODE_CREATE != 0 && comm.rank() == 0 {
71            let mut oo = OpenOptions::new();
72            oo.write(true);
73            if mode & MODE_EXCL != 0 {
74                oo.create_new(true);
75            } else {
76                oo.create(true);
77            }
78            let _ = oo.open(&path)?;
79        }
80        comm.barrier();
81
82        let mut oo = OpenOptions::new();
83        let read = mode & (MODE_RDONLY | MODE_RDWR) != 0;
84        let write = mode & (MODE_WRONLY | MODE_RDWR | MODE_APPEND | MODE_CREATE) != 0;
85        oo.read(read || !write).write(write);
86        if mode & MODE_APPEND != 0 {
87            oo.append(true);
88        }
89        let inner = oo.open(&path)?;
90
91        Ok(File {
92            inner,
93            comm: comm.duplicate(),
94            path,
95            delete_on_close: mode & MODE_DELETE_ON_CLOSE != 0,
96        })
97    }
98
99    /// Write `data` at byte `offset` (`MPI_File_write_at`).
100    pub fn write_at<T: Equivalence>(&mut self, offset: u64, data: &[T]) -> io::Result<()> {
101        self.inner.seek(SeekFrom::Start(offset))?;
102        self.inner.write_all(as_bytes(data))
103    }
104
105    /// Read into `buf` from byte `offset` (`MPI_File_read_at`).
106    pub fn read_at<T: Equivalence>(&mut self, offset: u64, buf: &mut [T]) -> io::Result<()> {
107        self.inner.seek(SeekFrom::Start(offset))?;
108        self.inner.read_exact(as_bytes_mut(buf))
109    }
110
111    /// Collective write at an explicit offset (`MPI_File_write_at_all`).
112    pub fn write_at_all<T: Equivalence>(&mut self, offset: u64, data: &[T]) -> io::Result<()> {
113        self.comm.barrier();
114        let r = self.write_at(offset, data);
115        self.inner.flush()?;
116        self.comm.barrier();
117        r
118    }
119
120    /// Collective read at an explicit offset (`MPI_File_read_at_all`).
121    pub fn read_at_all<T: Equivalence>(&mut self, offset: u64, buf: &mut [T]) -> io::Result<()> {
122        self.comm.barrier();
123        let r = self.read_at(offset, buf);
124        self.comm.barrier();
125        r
126    }
127
128    /// The current file size in bytes (`MPI_File_get_size`).
129    pub fn size(&self) -> io::Result<u64> {
130        Ok(self.inner.metadata()?.len())
131    }
132
133    /// Set the file size in bytes (`MPI_File_set_size`).
134    pub fn set_size(&mut self, size: u64) -> io::Result<()> {
135        self.inner.set_len(size)
136    }
137
138    /// Flush buffered data to storage (`MPI_File_sync`).
139    pub fn sync(&mut self) -> io::Result<()> {
140        self.inner.sync_all()
141    }
142
143    /// Collectively close the file (`MPI_File_close`), deleting it if it was
144    /// opened with [`MODE_DELETE_ON_CLOSE`].
145    pub fn close(self) -> io::Result<()> {
146        // Drop does the work; this consumes `self` for an explicit close.
147        drop(self);
148        Ok(())
149    }
150}
151
152impl Drop for File {
153    fn drop(&mut self) {
154        let _ = self.inner.sync_all();
155        if self.delete_on_close {
156            self.comm.barrier();
157            if self.comm.rank() == 0 {
158                let _ = std::fs::remove_file(&self.path);
159            }
160        }
161    }
162}