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(super) fn check_json_size(
23 value: &impl Serialize,
24 max_bytes: u64,
25) -> Result<(), PersistenceError> {
26 let mut writer = ic_host_artifacts::artifact::BoundedWriter::new(io::sink(), max_bytes);
27 let result = serde_json::to_writer_pretty(&mut writer, value);
28 if writer.limit_exceeded() {
29 return Err(PersistenceError::RecordTooLarge { limit: max_bytes });
30 }
31 result?;
32 Ok(())
33}
34
35pub fn write_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
40where
41 T: Serialize,
42{
43 let bytes = serde_json::to_vec_pretty(value)?;
44 replace_bytes_at_barriers(path, &bytes, |_| {}).map_err(PersistenceError::from)
45}
46
47pub fn create_json_durable<T>(path: &Path, value: &T) -> Result<(), PersistenceError>
52where
53 T: Serialize,
54{
55 let bytes = serde_json::to_vec_pretty(value)?;
56 create_bytes_at_barriers(path, &bytes, || {}, || {}).map_err(PersistenceError::from)
57}
58
59pub fn read_json<T>(path: &Path, max_bytes: u64) -> Result<T, PersistenceError>
67where
68 T: DeserializeOwned,
69{
70 #[cfg(unix)]
71 {
72 let bytes = ic_host_fs::read::read_file_no_follow(
73 path,
74 usize::try_from(max_bytes).unwrap_or(usize::MAX),
75 )
76 .map_err(|error| record_read_error(error, max_bytes))?;
77 Ok(serde_json::from_slice(&bytes)?)
78 }
79 #[cfg(not(unix))]
80 {
81 let _ = (path, max_bytes);
82 Err(io::Error::from(io::ErrorKind::Unsupported).into())
83 }
84}
85
86#[cfg(unix)]
87fn record_read_error(
88 error: ic_host_artifacts::artifact::ArtifactError,
89 max_bytes: u64,
90) -> PersistenceError {
91 use ic_host_artifacts::artifact::ArtifactError;
92 match error {
93 ArtifactError::NotRegularFile => PersistenceError::Io(io::Error::new(
94 io::ErrorKind::InvalidInput,
95 "record must be a regular file",
96 )),
97 ArtifactError::LimitExceeded { .. } => {
98 PersistenceError::RecordTooLarge { limit: max_bytes }
99 }
100 error => PersistenceError::Io(error.into()),
101 }
102}
103
104fn create_bytes_at_barriers(
105 path: &Path,
106 bytes: &[u8],
107 mut before_publication: impl FnMut(),
108 mut after_directory_sync: impl FnMut(),
109) -> io::Result<()> {
110 let parent = path
111 .parent()
112 .filter(|parent| !parent.as_os_str().is_empty())
113 .unwrap_or_else(|| Path::new("."));
114 create_private_parents(parent)?;
115
116 let (temp_path, mut temp_file) = create_sibling_temp(path, parent)?;
117 if let Err(error) = temp_file
118 .write_all(bytes)
119 .and_then(|()| temp_file.sync_all())
120 {
121 drop(temp_file);
122 let _ = fs::remove_file(&temp_path);
123 return Err(error);
124 }
125 drop(temp_file);
126 before_publication();
127
128 if let Err(error) = fs::hard_link(&temp_path, path) {
129 let _ = fs::remove_file(&temp_path);
130 return Err(error);
131 }
132 fs::remove_file(&temp_path)?;
133 File::open(parent)?.sync_all()?;
134 after_directory_sync();
135 Ok(())
136}
137
138#[derive(Clone, Copy, Debug, Eq, PartialEq)]
139pub(crate) enum DurableWriteBarrier {
140 BeforeRename,
141 AfterDirectorySync,
142}
143
144fn replace_bytes_at_barriers(
145 path: &Path,
146 bytes: &[u8],
147 mut barrier: impl FnMut(DurableWriteBarrier),
148) -> io::Result<()> {
149 let parent = path
150 .parent()
151 .filter(|parent| !parent.as_os_str().is_empty())
152 .unwrap_or_else(|| Path::new("."));
153 create_private_parents(parent)?;
154
155 let (temp_path, mut temp_file) = create_sibling_temp(path, parent)?;
156 if let Err(error) = temp_file
157 .write_all(bytes)
158 .and_then(|()| temp_file.sync_all())
159 {
160 drop(temp_file);
161 let _ = fs::remove_file(&temp_path);
162 return Err(error);
163 }
164 drop(temp_file);
165 barrier(DurableWriteBarrier::BeforeRename);
166
167 if let Err(error) = fs::rename(&temp_path, path) {
168 let _ = fs::remove_file(&temp_path);
169 return Err(error);
170 }
171
172 File::open(parent)?.sync_all()?;
173 barrier(DurableWriteBarrier::AfterDirectorySync);
174 Ok(())
175}
176
177#[cfg(test)]
178pub(crate) fn write_json_durable_at_barriers<T>(
179 path: &Path,
180 value: &T,
181 barrier: impl FnMut(DurableWriteBarrier),
182) -> Result<(), PersistenceError>
183where
184 T: Serialize,
185{
186 let bytes = serde_json::to_vec_pretty(value)?;
187 replace_bytes_at_barriers(path, &bytes, barrier).map_err(PersistenceError::from)
188}
189
190#[cfg(test)]
191pub(crate) fn create_json_durable_at_barriers<T>(
192 path: &Path,
193 value: &T,
194 before_publication: impl FnMut(),
195 after_directory_sync: impl FnMut(),
196) -> Result<(), PersistenceError>
197where
198 T: Serialize,
199{
200 let bytes = serde_json::to_vec_pretty(value)?;
201 create_bytes_at_barriers(path, &bytes, before_publication, after_directory_sync)
202 .map_err(PersistenceError::from)
203}
204
205fn create_sibling_temp(path: &Path, parent: &Path) -> io::Result<(PathBuf, File)> {
206 let file_name = path.file_name().ok_or_else(|| {
207 io::Error::new(
208 io::ErrorKind::InvalidInput,
209 format!("durable write target has no file name: {}", path.display()),
210 )
211 })?;
212
213 for _ in 0..64 {
214 let sequence = TEMP_SEQUENCE.fetch_add(1, Ordering::Relaxed);
215 let mut temp_name = OsString::from(".");
216 temp_name.push(file_name);
217 temp_name.push(format!(".ic-backup-tmp-{}-{sequence}", std::process::id()));
218 let temp_path = parent.join(temp_name);
219 let mut options = OpenOptions::new();
220 options.write(true).create_new(true);
221 #[cfg(unix)]
222 {
223 use std::os::unix::fs::OpenOptionsExt;
224 options.mode(0o600);
225 }
226 match options.open(&temp_path) {
227 Ok(file) => return Ok((temp_path, file)),
228 Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
229 Err(error) => return Err(error),
230 }
231 }
232
233 Err(io::Error::new(
234 io::ErrorKind::AlreadyExists,
235 format!(
236 "could not allocate a unique sibling temporary file for {}",
237 path.display()
238 ),
239 ))
240}
241
242#[cfg(test)]
247mod regressions;
248#[cfg(test)]
249mod tests;
250
251fn create_private_parents(parent: &Path) -> io::Result<()> {
252 let mut missing = Vec::new();
253 let mut current = parent;
254 while !current.try_exists()? {
255 missing.push(current.to_path_buf());
256 current = current
257 .parent()
258 .filter(|path| !path.as_os_str().is_empty())
259 .unwrap_or_else(|| Path::new("."));
260 }
261 let mut builder = fs::DirBuilder::new();
262 builder.recursive(true);
263 #[cfg(unix)]
264 {
265 use std::os::unix::fs::DirBuilderExt;
266 builder.mode(0o700);
267 }
268 builder.create(parent)?;
269 for directory in missing {
271 File::open(&directory)?.sync_all()?;
272 let ancestor = directory
273 .parent()
274 .filter(|path| !path.as_os_str().is_empty())
275 .unwrap_or_else(|| Path::new("."));
276 File::open(ancestor)?.sync_all()?;
277 }
278 Ok(())
279}