ic_backup/ops/persistence/command_lifetime_lock/
mod.rs1use super::file_lock::{self, FileLockError};
4use crate::model::command_custody::{CommandCustodyRecord, CommandCustodyRecordError};
5use std::{
6 fs::{self, File},
7 io,
8 path::{Path, PathBuf},
9 process::{Child, Command},
10 thread,
11 time::{Duration, Instant},
12};
13use thiserror::Error;
14
15const QUIESCENCE_GRACE: Duration = Duration::from_millis(250);
16pub const COMMAND_CUSTODY_DESCRIPTOR_ENV: &str = "IC_BACKUP_COMMAND_CUSTODY_FD";
18
19#[derive(Debug)]
24pub struct CommandLifetimeLock {
25 file: File,
26 path: PathBuf,
27 record: CommandCustodyRecord,
28 dispatched: bool,
29}
30
31impl CommandLifetimeLock {
32 pub fn acquire(
40 journal: &Path,
41 operation_sequence: u64,
42 ) -> Result<Self, CommandLifetimeLockError> {
43 let journal = resolve_journal(journal)?;
44 let path = lock_path(&journal, operation_sequence);
45 let file = file_lock::acquire(&path).map_err(|error| project_error(&path, error))?;
46 let (device, inode) = file_identity(&file.metadata()?)?;
47 let record = CommandCustodyRecord::new(journal, operation_sequence, device, inode)?;
48 Ok(Self {
49 file,
50 path,
51 record,
52 dispatched: false,
53 })
54 }
55
56 #[must_use]
58 pub fn record(&self) -> &CommandCustodyRecord {
59 &self.record
60 }
61
62 #[must_use]
64 pub fn path(&self) -> &Path {
65 &self.path
66 }
67
68 pub fn spawn(&mut self, mut command: Command) -> Result<Child, CommandLifetimeLockError> {
79 if self.dispatched {
80 return Err(CommandLifetimeLockError::AlreadyDispatched {
81 path: self.path.clone(),
82 });
83 }
84 require_identity(&self.record, &fs::symlink_metadata(&self.path)?, &self.path)?;
85 self.dispatched = true;
86 #[cfg(unix)]
87 {
88 use command_fds::CommandFdExt;
89 use std::os::fd::AsRawFd;
90 let descriptor = rustix::io::fcntl_dupfd_cloexec(&self.file, 3)
91 .map_err(|error| io::Error::from_raw_os_error(error.raw_os_error()))?;
92 command.env(
93 COMMAND_CUSTODY_DESCRIPTOR_ENV,
94 descriptor.as_raw_fd().to_string(),
95 );
96 command.preserved_fds(vec![descriptor]);
97 let child = command.spawn();
98 drop(command);
99 child.map_err(CommandLifetimeLockError::from)
100 }
101 #[cfg(not(unix))]
102 {
103 let _ = command;
104 Err(io::Error::from(io::ErrorKind::Unsupported).into())
105 }
106 }
107
108 pub fn finish(self) -> Result<CommandQuiescenceGuard, CommandLifetimeLockError> {
117 let Self { file, record, .. } = self;
118 drop(file);
119 let deadline = Instant::now() + QUIESCENCE_GRACE;
120 loop {
121 match CommandQuiescenceGuard::acquire(&record) {
122 Err(CommandLifetimeLockError::InFlight { .. }) if Instant::now() < deadline => {
123 thread::sleep(Duration::from_millis(5));
124 }
125 result => return result,
126 }
127 }
128 }
129}
130
131#[derive(Debug)]
136pub struct CommandQuiescenceGuard {
137 _file: File,
138 record: CommandCustodyRecord,
139}
140
141impl CommandQuiescenceGuard {
142 pub fn acquire(expected: &CommandCustodyRecord) -> Result<Self, CommandLifetimeLockError> {
147 let path = lock_path(expected.journal(), expected.operation_sequence());
148 let file =
149 file_lock::acquire_existing(&path).map_err(|error| project_error(&path, error))?;
150 require_identity(expected, &file.metadata()?, &path)?;
151 Ok(Self {
152 _file: file,
153 record: expected.clone(),
154 })
155 }
156
157 #[must_use]
159 pub fn record(&self) -> &CommandCustodyRecord {
160 &self.record
161 }
162}
163
164fn resolve_journal(path: &Path) -> Result<PathBuf, CommandLifetimeLockError> {
165 let name = path
166 .file_name()
167 .ok_or_else(|| CommandLifetimeLockError::InvalidJournal {
168 path: path.to_path_buf(),
169 })?;
170 let parent = path
171 .parent()
172 .filter(|parent| !parent.as_os_str().is_empty())
173 .unwrap_or_else(|| Path::new("."));
174 let parent = parent.canonicalize()?;
175 if !parent.is_dir() {
176 return Err(io::Error::from(io::ErrorKind::NotADirectory).into());
177 }
178 let journal = parent.join(name);
179 match fs::symlink_metadata(&journal) {
180 Ok(metadata) if !metadata.is_file() => {
181 return Err(CommandLifetimeLockError::InvalidJournal { path: journal });
182 }
183 Err(error) if error.kind() != io::ErrorKind::NotFound => return Err(error.into()),
184 _ => {}
185 }
186 Ok(journal)
187}
188
189fn lock_path(journal: &Path, operation_sequence: u64) -> PathBuf {
190 let mut path = journal.as_os_str().to_os_string();
191 path.push(format!(".command-{operation_sequence}.lock"));
192 PathBuf::from(path)
193}
194
195fn file_identity(metadata: &fs::Metadata) -> io::Result<(u64, u64)> {
196 if !metadata.is_file() {
197 return Err(io::Error::from(io::ErrorKind::InvalidInput));
198 }
199 #[cfg(unix)]
200 {
201 use std::os::unix::fs::MetadataExt;
202 Ok((metadata.dev(), metadata.ino()))
203 }
204 #[cfg(not(unix))]
205 {
206 let _ = metadata;
207 Err(io::Error::from(io::ErrorKind::Unsupported))
208 }
209}
210
211fn require_identity(
212 expected: &CommandCustodyRecord,
213 metadata: &fs::Metadata,
214 path: &Path,
215) -> Result<(), CommandLifetimeLockError> {
216 if !metadata.is_file() {
217 return Err(CommandLifetimeLockError::IdentityChanged {
218 path: path.to_path_buf(),
219 });
220 }
221 let (device, inode) = file_identity(metadata)?;
222 if expected.matches_file(device, inode) {
223 return Ok(());
224 }
225 Err(CommandLifetimeLockError::IdentityChanged {
226 path: path.to_path_buf(),
227 })
228}
229
230fn project_error(path: &Path, error: FileLockError) -> CommandLifetimeLockError {
231 match error {
232 FileLockError::Locked => CommandLifetimeLockError::InFlight {
233 path: path.to_path_buf(),
234 },
235 FileLockError::UnsafeEntry { kind } => CommandLifetimeLockError::UnsafeEntry {
236 path: path.to_path_buf(),
237 kind,
238 },
239 FileLockError::Io(error) => CommandLifetimeLockError::Io(error),
240 }
241}
242
243#[derive(Debug, Error)]
245pub enum CommandLifetimeLockError {
246 #[error("command custody remains in flight: {path:?}")]
248 InFlight {
249 path: PathBuf,
251 },
252 #[error("command custody identity changed: {path:?}")]
254 IdentityChanged {
255 path: PathBuf,
257 },
258 #[error("command spawn already attempted: {path:?}")]
260 AlreadyDispatched {
261 path: PathBuf,
263 },
264 #[error("unsafe command custody entry at {path:?}: {kind}")]
266 UnsafeEntry {
267 path: PathBuf,
269 kind: String,
271 },
272 #[error("invalid command journal location: {path:?}")]
274 InvalidJournal {
275 path: PathBuf,
277 },
278 #[error(transparent)]
280 Record(#[from] CommandCustodyRecordError),
281 #[error(transparent)]
283 Io(#[from] io::Error),
284}
285
286#[cfg(all(test, unix))]
287mod tests;