Skip to main content

ic_backup/ops/persistence/json/
mod.rs

1//! Module: `persistence::json`
2//!
3//! Responsibility: read and durably create or replace JSON persistence documents.
4//! Does not own: document validation, layout paths, or integrity checks.
5//! Boundary: provides explicit create-only and replace filesystem primitives.
6
7use crate::ops::persistence::PersistenceError;
8
9use std::{
10    ffi::OsString,
11    fs::{self, File, OpenOptions},
12    io::{self, Write},
13    path::{Path, PathBuf},
14    sync::atomic::{AtomicU64, Ordering},
15};
16
17use serde::{Serialize, de::DeserializeOwned};
18#[cfg(unix)]
19use std::io::Read;
20
21static TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(0);
22
23/// Durably replace a machine record using a sibling temporary and rename.
24///
25/// # Errors
26/// Returns encoding or IO failures; a lost post-rename response requires reconciliation.
27pub fn write_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
28where
29    T: Serialize,
30{
31    let bytes = serde_json::to_vec_pretty(value)?;
32    replace_bytes_at_barriers(path, &bytes, |_| {}).map_err(PersistenceError::from)
33}
34
35/// Publish a new machine record without replacing an existing entry.
36///
37/// # Errors
38/// Returns encoding, existing-destination or IO failures.
39pub fn create_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
40where
41    T: Serialize,
42{
43    let bytes = serde_json::to_vec_pretty(value)?;
44    create_bytes_at_barriers(path, &bytes, || {}, || {}).map_err(PersistenceError::from)
45}
46
47/// Read one regular no-follow JSON file within an explicit byte limit.
48///
49/// # Errors
50/// Rejects unsafe files, excessive bytes, invalid JSON and filesystem failures.
51pub fn read_json<T>(path: &Path, max_bytes: u64) -> Result<T, PersistenceError>
52where
53    T: DeserializeOwned,
54{
55    #[cfg(unix)]
56    {
57        let file = {
58            use rustix::fs::{FileType, Mode, OFlags, fstat, open};
59            let fd = open(
60                path,
61                OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::NONBLOCK | OFlags::CLOEXEC,
62                Mode::empty(),
63            )
64            .map_err(|error| io::Error::from_raw_os_error(error.raw_os_error()))?;
65            let metadata =
66                fstat(&fd).map_err(|error| io::Error::from_raw_os_error(error.raw_os_error()))?;
67            if !FileType::from_raw_mode(metadata.st_mode).is_file() {
68                return Err(io::Error::new(
69                    io::ErrorKind::InvalidInput,
70                    "record must be a regular file",
71                )
72                .into());
73            }
74            File::from(fd)
75        };
76        let mut bytes = Vec::new();
77        file.take(max_bytes.saturating_add(1))
78            .read_to_end(&mut bytes)?;
79        if bytes.len() > usize::try_from(max_bytes).unwrap_or(usize::MAX) {
80            return Err(PersistenceError::RecordTooLarge { limit: max_bytes });
81        }
82        Ok(serde_json::from_slice(&bytes)?)
83    }
84    #[cfg(not(unix))]
85    {
86        let _ = (path, max_bytes);
87        Err(io::Error::from(io::ErrorKind::Unsupported).into())
88    }
89}
90
91fn create_bytes_at_barriers(
92    path: &Path,
93    bytes: &[u8],
94    mut before_publication: impl FnMut(),
95    mut after_directory_sync: impl FnMut(),
96) -> io::Result<()> {
97    let parent = path
98        .parent()
99        .filter(|parent| !parent.as_os_str().is_empty())
100        .unwrap_or_else(|| Path::new("."));
101    create_private_parents(parent)?;
102
103    let (temp_path, mut temp_file) = create_sibling_temp(path, parent)?;
104    if let Err(error) = temp_file
105        .write_all(bytes)
106        .and_then(|()| temp_file.sync_all())
107    {
108        drop(temp_file);
109        let _ = fs::remove_file(&temp_path);
110        return Err(error);
111    }
112    drop(temp_file);
113    before_publication();
114
115    if let Err(error) = fs::hard_link(&temp_path, path) {
116        let _ = fs::remove_file(&temp_path);
117        return Err(error);
118    }
119    fs::remove_file(&temp_path)?;
120    File::open(parent)?.sync_all()?;
121    after_directory_sync();
122    Ok(())
123}
124
125#[derive(Clone, Copy, Debug, Eq, PartialEq)]
126pub(crate) enum DurableWriteBarrier {
127    BeforeRename,
128    AfterDirectorySync,
129}
130
131fn replace_bytes_at_barriers(
132    path: &Path,
133    bytes: &[u8],
134    mut barrier: impl FnMut(DurableWriteBarrier),
135) -> io::Result<()> {
136    let parent = path
137        .parent()
138        .filter(|parent| !parent.as_os_str().is_empty())
139        .unwrap_or_else(|| Path::new("."));
140    create_private_parents(parent)?;
141
142    let (temp_path, mut temp_file) = create_sibling_temp(path, parent)?;
143    if let Err(error) = temp_file
144        .write_all(bytes)
145        .and_then(|()| temp_file.sync_all())
146    {
147        drop(temp_file);
148        let _ = fs::remove_file(&temp_path);
149        return Err(error);
150    }
151    drop(temp_file);
152    barrier(DurableWriteBarrier::BeforeRename);
153
154    if let Err(error) = fs::rename(&temp_path, path) {
155        let _ = fs::remove_file(&temp_path);
156        return Err(error);
157    }
158
159    File::open(parent)?.sync_all()?;
160    barrier(DurableWriteBarrier::AfterDirectorySync);
161    Ok(())
162}
163
164#[cfg(test)]
165pub(crate) fn write_json_durable_at_barriers<T>(
166    path: &Path,
167    value: &T,
168    barrier: impl FnMut(DurableWriteBarrier),
169) -> Result<(), PersistenceError>
170where
171    T: Serialize,
172{
173    let bytes = serde_json::to_vec_pretty(value)?;
174    replace_bytes_at_barriers(path, &bytes, barrier).map_err(PersistenceError::from)
175}
176
177#[cfg(test)]
178pub(crate) fn create_json_durable_at_barriers<T>(
179    path: &Path,
180    value: &T,
181    before_publication: impl FnMut(),
182    after_directory_sync: impl FnMut(),
183) -> Result<(), PersistenceError>
184where
185    T: Serialize,
186{
187    let bytes = serde_json::to_vec_pretty(value)?;
188    create_bytes_at_barriers(path, &bytes, before_publication, after_directory_sync)
189        .map_err(PersistenceError::from)
190}
191
192fn create_sibling_temp(path: &Path, parent: &Path) -> io::Result<(PathBuf, File)> {
193    let file_name = path.file_name().ok_or_else(|| {
194        io::Error::new(
195            io::ErrorKind::InvalidInput,
196            format!("durable write target has no file name: {}", path.display()),
197        )
198    })?;
199
200    for _ in 0..64 {
201        let sequence = TEMP_SEQUENCE.fetch_add(1, Ordering::Relaxed);
202        let mut temp_name = OsString::from(".");
203        temp_name.push(file_name);
204        temp_name.push(format!(".ic-backup-tmp-{}-{sequence}", std::process::id()));
205        let temp_path = parent.join(temp_name);
206        let mut options = OpenOptions::new();
207        options.write(true).create_new(true);
208        #[cfg(unix)]
209        {
210            use std::os::unix::fs::OpenOptionsExt;
211            options.mode(0o600);
212        }
213        match options.open(&temp_path) {
214            Ok(file) => return Ok((temp_path, file)),
215            Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
216            Err(error) => return Err(error),
217        }
218    }
219
220    Err(io::Error::new(
221        io::ErrorKind::AlreadyExists,
222        format!(
223            "could not allocate a unique sibling temporary file for {}",
224            path.display()
225        ),
226    ))
227}
228
229// -----------------------------------------------------------------------------
230// Tests
231// -----------------------------------------------------------------------------
232
233#[cfg(test)]
234mod regressions;
235#[cfg(test)]
236mod tests;
237
238fn create_private_parents(parent: &Path) -> io::Result<()> {
239    let mut missing = Vec::new();
240    let mut current = parent;
241    while !current.try_exists()? {
242        missing.push(current.to_path_buf());
243        current = current
244            .parent()
245            .filter(|path| !path.as_os_str().is_empty())
246            .unwrap_or_else(|| Path::new("."));
247    }
248    let mut builder = fs::DirBuilder::new();
249    builder.recursive(true);
250    #[cfg(unix)]
251    {
252        use std::os::unix::fs::DirBuilderExt;
253        builder.mode(0o700);
254    }
255    builder.create(parent)?;
256    // Persist each newly created directory and its link from the existing root.
257    for directory in missing {
258        File::open(&directory)?.sync_all()?;
259        let ancestor = directory
260            .parent()
261            .filter(|path| !path.as_os_str().is_empty())
262            .unwrap_or_else(|| Path::new("."));
263        File::open(ancestor)?.sync_all()?;
264    }
265    Ok(())
266}