ic_backup/ops/persistence/json/
mod.rs1use 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
23pub 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
35pub 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
47pub 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#[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 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}