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)]
138pub struct CommandQuiescenceGuard {
139 file: File,
140 record: CommandCustodyRecord,
141}
142
143impl CommandQuiescenceGuard {
144 pub fn acquire(expected: &CommandCustodyRecord) -> Result<Self, CommandLifetimeLockError> {
149 let path = lock_path(expected.journal(), expected.operation_sequence());
150 let file =
151 file_lock::acquire_existing(&path).map_err(|error| project_error(&path, error))?;
152 require_identity(expected, &file.metadata()?, &path)?;
153 Ok(Self {
154 file,
155 record: expected.clone(),
156 })
157 }
158
159 #[must_use]
161 pub fn record(&self) -> &CommandCustodyRecord {
162 &self.record
163 }
164}
165
166impl Drop for CommandQuiescenceGuard {
167 fn drop(&mut self) {
168 #[cfg(unix)]
172 file_lock::unlock(&self.file);
173 }
174}
175
176fn resolve_journal(path: &Path) -> Result<PathBuf, CommandLifetimeLockError> {
177 let name = path
178 .file_name()
179 .ok_or_else(|| CommandLifetimeLockError::InvalidJournal {
180 path: path.to_path_buf(),
181 })?;
182 let parent = path
183 .parent()
184 .filter(|parent| !parent.as_os_str().is_empty())
185 .unwrap_or_else(|| Path::new("."));
186 let parent = parent.canonicalize()?;
187 if !parent.is_dir() {
188 return Err(io::Error::from(io::ErrorKind::NotADirectory).into());
189 }
190 let journal = parent.join(name);
191 match fs::symlink_metadata(&journal) {
192 Ok(metadata) if !metadata.is_file() => {
193 return Err(CommandLifetimeLockError::InvalidJournal { path: journal });
194 }
195 Err(error) if error.kind() != io::ErrorKind::NotFound => return Err(error.into()),
196 _ => {}
197 }
198 Ok(journal)
199}
200
201fn lock_path(journal: &Path, operation_sequence: u64) -> PathBuf {
202 let mut path = journal.as_os_str().to_os_string();
203 path.push(format!(".command-{operation_sequence}.lock"));
204 PathBuf::from(path)
205}
206
207fn file_identity(metadata: &fs::Metadata) -> io::Result<(u64, u64)> {
208 if !metadata.is_file() {
209 return Err(io::Error::from(io::ErrorKind::InvalidInput));
210 }
211 #[cfg(unix)]
212 {
213 use std::os::unix::fs::MetadataExt;
214 Ok((metadata.dev(), metadata.ino()))
215 }
216 #[cfg(not(unix))]
217 {
218 let _ = metadata;
219 Err(io::Error::from(io::ErrorKind::Unsupported))
220 }
221}
222
223fn require_identity(
224 expected: &CommandCustodyRecord,
225 metadata: &fs::Metadata,
226 path: &Path,
227) -> Result<(), CommandLifetimeLockError> {
228 if !metadata.is_file() {
229 return Err(CommandLifetimeLockError::IdentityChanged {
230 path: path.to_path_buf(),
231 });
232 }
233 let (device, inode) = file_identity(metadata)?;
234 if expected.matches_file(device, inode) {
235 return Ok(());
236 }
237 Err(CommandLifetimeLockError::IdentityChanged {
238 path: path.to_path_buf(),
239 })
240}
241
242fn project_error(path: &Path, error: FileLockError) -> CommandLifetimeLockError {
243 match error {
244 FileLockError::Locked => CommandLifetimeLockError::InFlight {
245 path: path.to_path_buf(),
246 },
247 FileLockError::UnsafeEntry { kind } => CommandLifetimeLockError::UnsafeEntry {
248 path: path.to_path_buf(),
249 kind,
250 },
251 FileLockError::Io(error) => CommandLifetimeLockError::Io(error),
252 }
253}
254
255#[derive(Debug, Error)]
257pub enum CommandLifetimeLockError {
258 #[error("command custody remains in flight: {path:?}")]
260 InFlight {
261 path: PathBuf,
263 },
264 #[error("command custody identity changed: {path:?}")]
266 IdentityChanged {
267 path: PathBuf,
269 },
270 #[error("command spawn already attempted: {path:?}")]
272 AlreadyDispatched {
273 path: PathBuf,
275 },
276 #[error("unsafe command custody entry at {path:?}: {kind}")]
278 UnsafeEntry {
279 path: PathBuf,
281 kind: String,
283 },
284 #[error("invalid command journal location: {path:?}")]
286 InvalidJournal {
287 path: PathBuf,
289 },
290 #[error(transparent)]
292 Record(#[from] CommandCustodyRecordError),
293 #[error(transparent)]
295 Io(#[from] io::Error),
296}
297
298#[cfg(all(test, unix))]
299mod tests;