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
19static TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(0);
20
21pub fn write_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
26where
27 T: Serialize,
28{
29 let bytes = serde_json::to_vec_pretty(value)?;
30 replace_bytes_at_barriers(path, &bytes, |_| {}).map_err(PersistenceError::from)
31}
32
33pub fn create_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
38where
39 T: Serialize,
40{
41 let bytes = serde_json::to_vec_pretty(value)?;
42 create_bytes_at_barriers(path, &bytes, || {}, || {}).map_err(PersistenceError::from)
43}
44
45pub fn read_json<T>(path: &Path, max_bytes: u64) -> Result<T, PersistenceError>
53where
54 T: DeserializeOwned,
55{
56 #[cfg(unix)]
57 {
58 let bytes = ic_host_tools::artifact::read_file_no_follow(
59 path,
60 usize::try_from(max_bytes).unwrap_or(usize::MAX),
61 )
62 .map_err(|error| record_read_error(error, max_bytes))?;
63 Ok(serde_json::from_slice(&bytes)?)
64 }
65 #[cfg(not(unix))]
66 {
67 let _ = (path, max_bytes);
68 Err(io::Error::from(io::ErrorKind::Unsupported).into())
69 }
70}
71
72#[cfg(unix)]
73fn record_read_error(
74 error: ic_host_tools::artifact::ArtifactError,
75 max_bytes: u64,
76) -> PersistenceError {
77 use ic_host_tools::artifact::ArtifactError;
78 match error {
79 ArtifactError::Io(error) => PersistenceError::Io(error),
80 ArtifactError::NotRegularFile => PersistenceError::Io(io::Error::new(
81 io::ErrorKind::InvalidInput,
82 "record must be a regular file",
83 )),
84 ArtifactError::LimitExceeded { .. } => {
85 PersistenceError::RecordTooLarge { limit: max_bytes }
86 }
87 error => PersistenceError::Io(io::Error::other(error)),
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}