1use std::fs::{File, OpenOptions};
8use std::io::{Read, Seek, Write};
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicU64, Ordering};
11use std::sync::{Mutex, OnceLock};
12
13use serde::{Deserialize, Serialize};
14use sha2::{Digest, Sha256};
15
16const MAX_MODEL_ID_BYTES: usize = 4 * 1024;
17const MAX_STATE_BYTES: u64 = 64 * 1024;
18const QUARANTINE_VANISHED_REASON: &str = "quarantine vanished before deletion; the removal journal is retained, so retrying the removal resumes and completes it";
25
26#[derive(Debug, thiserror::Error)]
27pub enum ModelManagementError {
28 #[error("model-management I/O failed at {path}: {source}")]
29 Io {
30 path: PathBuf,
31 #[source]
32 source: std::io::Error,
33 },
34 #[error("invalid model-management state at {path}: {message}")]
35 InvalidState { path: PathBuf, message: String },
36 #[error("model-management identity mismatch at {path}: expected {expected}, found {actual}")]
37 IdentityMismatch {
38 path: PathBuf,
39 expected: String,
40 actual: String,
41 },
42 #[error("no CAR install receipt exists for {model_id}")]
43 MissingReceipt { model_id: String },
44 #[error("unsafe CAR-managed path for {model_id}: {path} ({reason})")]
45 UnsafeManagedPath {
46 model_id: String,
47 path: PathBuf,
48 reason: String,
49 },
50 #[error("local model {model_id} is in use")]
51 ModelInUse { model_id: String },
52 #[error("local model-management operation is already active for {model_id}")]
53 MutationInProgress { model_id: String },
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(rename_all = "snake_case")]
58pub enum ManagedArtifactKind {
59 Symlink,
60 Directory,
61 File,
62}
63
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(deny_unknown_fields)]
66pub struct InstallReceipt {
67 pub model_id: String,
68 pub managed_path: PathBuf,
69 pub artifact_kind: ManagedArtifactKind,
70 pub source_model_id: String,
71 pub source_revision: Option<String>,
72 pub creation_generation: u64,
73 pub shared_cache_references: Vec<PathBuf>,
74 pub adopted: bool,
75}
76
77pub const USAGE_STAMP_INTERVAL: std::time::Duration = std::time::Duration::from_secs(10 * 60);
79
80const USAGE_TRACKING_MARKER: &str = "tracking-since.json";
81const USAGE_COVERAGE: &str = "coverage.json";
82
83pub const COVERAGE_GAP_SECS: u64 = 60 * 60;
87
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
92struct UsageCoverage {
93 run_start: u64,
94 last_beat: u64,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
100pub struct UsageStamp {
101 pub model_id: String,
102 pub last_used: u64,
103}
104
105#[derive(Debug, Serialize, Deserialize)]
106struct UsageTracking {
107 since: u64,
108}
109
110fn unix_now() -> u64 {
111 std::time::SystemTime::now()
112 .duration_since(std::time::UNIX_EPOCH)
113 .map(|d| d.as_secs())
114 .unwrap_or(0)
115}
116
117#[derive(Debug, Clone)]
118pub struct ModelManagementStore {
119 models_dir: PathBuf,
120 management_state_dir: PathBuf,
121 receipts_root: PathBuf,
122 tombstones_root: PathBuf,
123 installs_root: PathBuf,
124 removals_root: PathBuf,
125 leases_root: PathBuf,
126 mutation_locks_root: PathBuf,
127 usage_root: PathBuf,
131 usage_throttle: std::sync::Arc<Mutex<std::collections::HashMap<String, std::time::Instant>>>,
134 last_beat: std::sync::Arc<Mutex<Option<(u64, std::time::Instant)>>>,
138 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
142 removal_hook: Option<RemovalTestHook>,
143 #[cfg(test)]
144 identity_init_hook: Option<fn(&Path)>,
145}
146
147pub(crate) const fn directory_removal_supported() -> bool {
152 cfg!(any(target_os = "macos", target_os = "linux"))
153}
154
155#[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
156#[derive(Clone, Copy, Debug, PartialEq, Eq)]
157pub(crate) enum RemovalPhase {
158 AfterRename,
160 BeforeRootOpen,
163 AfterCapture,
166}
167
168#[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
169#[derive(Clone)]
170pub(crate) struct RemovalTestHook(std::sync::Arc<dyn Fn(RemovalPhase, &Path) + Send + Sync>);
171
172#[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
173impl std::fmt::Debug for RemovalTestHook {
174 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
175 f.write_str("RemovalTestHook")
176 }
177}
178
179impl ModelManagementStore {
180 pub fn new(state_root: PathBuf, models_dir: PathBuf) -> Self {
181 let state_root = crate::resource_policy::normalized_state_root_key(&state_root);
182 let management_state_dir = state_root.join("model-management");
183 let models_dir = crate::resource_policy::normalized_state_root_key(&models_dir);
187 let shared_coordination_root = models_dir
188 .parent()
189 .unwrap_or(&models_dir)
190 .join(".car-model-management");
191 let store = Self {
192 models_dir,
193 receipts_root: receipts_dir(&state_root),
194 tombstones_root: management_state_dir.join("tombstones"),
195 installs_root: management_state_dir.join("installs"),
196 removals_root: management_state_dir.join("removals"),
197 leases_root: shared_coordination_root.join("activity"),
198 mutation_locks_root: shared_coordination_root.join("mutation"),
199 usage_root: shared_coordination_root.join("usage"),
200 usage_throttle: Default::default(),
201 last_beat: Default::default(),
202 management_state_dir,
203 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
204 removal_hook: None,
205 #[cfg(test)]
206 identity_init_hook: None,
207 };
208 store.resume_pending_quarantines();
212 store
213 }
214
215 pub fn models_dir(&self) -> &Path {
216 &self.models_dir
217 }
218
219 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
220 pub(crate) fn with_removal_hook(
221 mut self,
222 hook: impl Fn(RemovalPhase, &Path) + Send + Sync + 'static,
223 ) -> Self {
224 self.removal_hook = Some(RemovalTestHook(std::sync::Arc::new(hook)));
225 self
226 }
227
228 pub fn management_state_dir(&self) -> &Path {
229 &self.management_state_dir
230 }
231
232 pub fn receipts_root(&self) -> &Path {
233 &self.receipts_root
234 }
235
236 pub fn receipt_path(&self, model_id: &str) -> PathBuf {
237 self.receipts_root.join(hashed_filename(model_id))
238 }
239
240 pub fn tombstones_root(&self) -> &Path {
241 &self.tombstones_root
242 }
243
244 pub fn tombstone_path(&self, model_id: &str) -> PathBuf {
245 self.tombstones_root.join(hashed_filename(model_id))
246 }
247
248 fn removal_journal_path(&self, model_id: &str) -> PathBuf {
249 self.removals_root.join(hashed_filename(model_id))
250 }
251
252 fn install_journal_path(&self, model_id: &str) -> PathBuf {
253 self.installs_root.join(hashed_filename(model_id))
254 }
255
256 pub fn leases_root(&self) -> &Path {
257 &self.leases_root
258 }
259
260 pub fn lease_path(&self, model_id: &str) -> PathBuf {
261 self.leases_root.join(hashed_filename(model_id))
262 }
263
264 pub fn mutation_lock_path(&self, model_id: &str) -> PathBuf {
265 self.mutation_locks_root.join(hashed_filename(model_id))
266 }
267
268 pub fn load_receipt(
269 &self,
270 model_id: &str,
271 ) -> Result<Option<InstallReceipt>, ModelManagementError> {
272 validate_model_id(model_id)?;
273 read_receipt_at(&self.receipt_path(model_id), model_id)
274 }
275
276 pub(crate) fn write_receipt(
277 &self,
278 receipt: &InstallReceipt,
279 ) -> Result<(), ModelManagementError> {
280 validate_model_id(&receipt.model_id)?;
281 if let Some(existing) = self.load_receipt(&receipt.model_id)? {
282 ensure_identity(
283 &self.receipt_path(&receipt.model_id),
284 &receipt.model_id,
285 &existing.model_id,
286 )?;
287 }
288 write_private_json(&self.receipt_path(&receipt.model_id), receipt)
289 }
290
291 pub fn can_remove(&self, model_id: &str) -> Result<bool, ModelManagementError> {
294 let Some(receipt) = self.load_receipt(model_id)? else {
295 return Ok(false);
296 };
297 self.validate_managed_artifact(&receipt)?;
298 Ok(
303 receipt.artifact_kind != ManagedArtifactKind::Directory
304 || directory_removal_supported(),
305 )
306 }
307
308 pub fn acquire_lease(&self, model_id: &str) -> Result<ModelLease, ModelManagementError> {
309 validate_model_id(model_id)?;
310 let mutation = self.open_identity_lock(&self.mutation_lock_path(model_id), model_id)?;
314 mutation
315 .lock_shared()
316 .map_err(|source| ModelManagementError::Io {
317 path: self.mutation_lock_path(model_id),
318 source,
319 })?;
320 let mut mutation = OwnedFileLock::new(mutation);
321 validate_lock_identity(&mut mutation, &self.mutation_lock_path(model_id), model_id)?;
322 let file = self.open_identity_lock(&self.lease_path(model_id), model_id)?;
323 file.lock_shared()
324 .map_err(|source| ModelManagementError::Io {
325 path: self.lease_path(model_id),
326 source,
327 })?;
328 let mut file = OwnedFileLock::new(file);
329 validate_lock_identity(&mut file, &self.lease_path(model_id), model_id)?;
330 drop(mutation);
331 self.note_used(model_id);
332 Ok(ModelLease { _file: file })
333 }
334
335 fn note_used(&self, model_id: &str) {
341 let now = std::time::Instant::now();
342 {
343 let mut written = self
344 .usage_throttle
345 .lock()
346 .unwrap_or_else(std::sync::PoisonError::into_inner);
347 if written
348 .get(model_id)
349 .is_some_and(|at| now.duration_since(*at) < USAGE_STAMP_INTERVAL)
350 {
351 return;
352 }
353 written.insert(model_id.to_string(), now);
354 }
355 let stamp = UsageStamp {
356 model_id: model_id.to_string(),
357 last_used: unix_now(),
358 };
359 if let Err(error) =
360 write_private_json(&self.usage_root.join(hashed_filename(model_id)), &stamp)
361 {
362 tracing::debug!(model = model_id, %error, "could not record model use");
363 }
364 }
365
366 pub fn start_usage_tracking(&self) {
376 let marker = self.usage_root.join(USAGE_TRACKING_MARKER);
377 if std::fs::symlink_metadata(&marker).is_ok() {
378 return;
379 }
380 if let Err(error) = write_private_json(&marker, &UsageTracking { since: unix_now() }) {
381 tracing::debug!(%error, "could not record when model usage tracking began");
382 }
383 }
384
385 #[cfg(test)]
388 pub(crate) fn set_usage_tracking_since_for_test(&self, since: u64) {
389 write_private_json(
390 &self.usage_root.join(USAGE_TRACKING_MARKER),
391 &UsageTracking { since },
392 )
393 .unwrap();
394 write_private_json(
395 &self.usage_root.join(USAGE_COVERAGE),
396 &UsageCoverage {
397 run_start: since,
398 last_beat: unix_now(),
399 },
400 )
401 .unwrap();
402 }
403
404 pub fn beat(&self) {
411 let now = unix_now();
412 let instant = std::time::Instant::now();
413 let path = self.usage_root.join(USAGE_COVERAGE);
414 let previous = read_private_state(&path)
415 .ok()
416 .and_then(|bytes| serde_json::from_slice::<UsageCoverage>(&bytes).ok());
417 let mut mine = self
418 .last_beat
419 .lock()
420 .unwrap_or_else(std::sync::PoisonError::into_inner);
421 let run_start = match (previous, *mine) {
422 (Some(prev), Some((wall, at))) if prev.last_beat == wall => {
423 if instant.duration_since(at).as_secs() <= COVERAGE_GAP_SECS {
424 prev.run_start
425 } else {
426 now
427 }
428 }
429 (Some(prev), _) if now.saturating_sub(prev.last_beat) <= COVERAGE_GAP_SECS => {
430 prev.run_start
431 }
432 _ => now,
433 };
434 let coverage = UsageCoverage {
435 run_start,
436 last_beat: now,
437 };
438 if let Err(error) = write_private_json(&path, &coverage) {
439 tracing::debug!(%error, "could not record model usage coverage");
440 }
441 *mine = Some((now, instant));
442 }
443
444 pub fn usage_covered_since(&self) -> Option<u64> {
446 read_private_state(&self.usage_root.join(USAGE_COVERAGE))
447 .ok()
448 .and_then(|bytes| serde_json::from_slice::<UsageCoverage>(&bytes).ok())
449 .map(|coverage| coverage.run_start)
450 }
451
452 pub fn usage_tracking_since(&self) -> Option<u64> {
454 read_private_state(&self.usage_root.join(USAGE_TRACKING_MARKER))
455 .ok()
456 .and_then(|bytes| serde_json::from_slice::<UsageTracking>(&bytes).ok())
457 .map(|marker| marker.since)
458 }
459
460 pub fn last_used(&self, model_id: &str) -> Option<u64> {
462 let path = self.usage_root.join(hashed_filename(model_id));
463 read_private_state(&path)
464 .ok()
465 .and_then(|bytes| serde_json::from_slice::<UsageStamp>(&bytes).ok())
466 .filter(|stamp| stamp.model_id == model_id)
467 .map(|stamp| stamp.last_used)
468 }
469
470 pub(crate) fn begin_mutation(
471 &self,
472 model_id: &str,
473 ) -> Result<ModelMutationGuard, ModelManagementError> {
474 validate_model_id(model_id)?;
475 let mut file = self.open_identity_lock(&self.mutation_lock_path(model_id), model_id)?;
476 validate_lock_identity(&mut file, &self.mutation_lock_path(model_id), model_id)?;
477 match file.try_lock() {
478 Ok(()) => Ok(ModelMutationGuard {
479 store: self.clone(),
480 model_id: model_id.to_string(),
481 _file: OwnedFileLock::new(file),
482 }),
483 Err(std::fs::TryLockError::WouldBlock) => {
484 Err(ModelManagementError::MutationInProgress {
485 model_id: model_id.to_string(),
486 })
487 }
488 Err(std::fs::TryLockError::Error(source)) => Err(ModelManagementError::Io {
489 path: self.mutation_lock_path(model_id),
490 source,
491 }),
492 }
493 }
494
495 pub fn model_in_use(&self, model_id: &str) -> Result<bool, ModelManagementError> {
496 validate_model_id(model_id)?;
497 let mut file = self.open_identity_lock(&self.lease_path(model_id), model_id)?;
498 validate_lock_identity(&mut file, &self.lease_path(model_id), model_id)?;
499 match file.try_lock() {
500 Ok(()) => {
501 let _lock = OwnedFileLock::new(file);
502 Ok(false)
503 }
504 Err(std::fs::TryLockError::WouldBlock) => Ok(true),
505 Err(std::fs::TryLockError::Error(source)) => Err(ModelManagementError::Io {
506 path: self.lease_path(model_id),
507 source,
508 }),
509 }
510 }
511
512 pub fn car_enabled(&self, model_id: &str) -> Result<bool, ModelManagementError> {
513 Ok(self.load_tombstone(model_id)?.is_none())
514 }
515
516 fn load_tombstone(
517 &self,
518 model_id: &str,
519 ) -> Result<Option<RemovalTombstone>, ModelManagementError> {
520 validate_model_id(model_id)?;
521 let path = self.tombstone_path(model_id);
522 let bytes = match read_private_state(&path) {
523 Ok(bytes) => bytes,
524 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
525 Err(source) => return Err(ModelManagementError::Io { path, source }),
526 };
527 let tombstone: RemovalTombstone =
528 serde_json::from_slice(&bytes).map_err(|error| ModelManagementError::InvalidState {
529 path: path.clone(),
530 message: error.to_string(),
531 })?;
532 ensure_identity(&path, model_id, &tombstone.model_id)?;
533 Ok(Some(tombstone))
534 }
535
536 pub(crate) fn clear_tombstone(&self, model_id: &str) -> Result<(), ModelManagementError> {
537 validate_model_id(model_id)?;
538 let _ = self.car_enabled(model_id)?;
539 let path = self.tombstone_path(model_id);
540 match std::fs::remove_file(&path) {
541 Ok(()) => sync_directory(path.parent().expect("tombstone path has parent"))
542 .map_err(|source| ModelManagementError::Io { path, source }),
543 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
544 Err(source) => Err(ModelManagementError::Io { path, source }),
545 }
546 }
547
548 pub(crate) fn record_managed_artifact(
552 &self,
553 model_id: &str,
554 source_model_id: &str,
555 source_revision: Option<String>,
556 creation_generation: u64,
557 adopted: bool,
558 managed_path: PathBuf,
559 ) -> Result<InstallReceipt, ModelManagementError> {
560 validate_model_id(model_id)?;
561 let artifact_kind = artifact_kind(&managed_path, model_id)?;
562 let shared_cache_references = collect_shared_references(&managed_path, model_id)?;
563 let receipt = InstallReceipt {
564 model_id: model_id.to_string(),
565 managed_path,
566 artifact_kind,
567 source_model_id: source_model_id.to_string(),
568 source_revision,
569 creation_generation,
570 shared_cache_references,
571 adopted,
572 };
573 self.validate_managed_artifact(&receipt)?;
574 self.write_receipt(&receipt)?;
575 self.clear_tombstone(model_id)?;
576 Ok(receipt)
577 }
578
579 pub(crate) fn begin_install_intent(
580 &self,
581 receipt: InstallReceipt,
582 staging_path: Option<&Path>,
583 ) -> Result<(), ModelManagementError> {
584 validate_model_id(&receipt.model_id)?;
585 validate_direct_child(&self.models_dir, &receipt.managed_path, &receipt.model_id)?;
586 let staging_path = staging_path.map(Path::to_path_buf);
587 let staging_identity = staging_path
588 .as_deref()
589 .map(|path| artifact_identity(path, &receipt.model_id))
590 .transpose()?;
591 if let Some(path) = staging_path.as_deref() {
592 validate_direct_child(&self.models_dir, path, &receipt.model_id)?;
593 }
594 let intent = InstallJournal {
595 model_id: receipt.model_id.clone(),
596 receipt,
597 staging_path,
598 staging_identity,
599 };
600 let path = self.install_journal_path(&intent.model_id);
601 if let Some(existing) = self.load_install_intent(&intent.model_id)? {
602 if existing == intent {
603 return Ok(());
604 }
605 return Err(ModelManagementError::InvalidState {
606 path,
607 message: "a different CAR install intent already exists for this model".into(),
608 });
609 }
610 if std::fs::symlink_metadata(&intent.receipt.managed_path).is_ok() {
611 return Err(ModelManagementError::UnsafeManagedPath {
612 model_id: intent.model_id,
613 path: intent.receipt.managed_path,
614 reason: "pre-existing unreceipted leaf cannot become CAR-owned through an install intent"
615 .into(),
616 });
617 }
618 write_private_json(&path, &intent)
619 }
620
621 pub(crate) fn install_receipt_for_publication(
622 &self,
623 model_id: &str,
624 source_model_id: &str,
625 creation_generation: u64,
626 adopted: bool,
627 managed_path: PathBuf,
628 artifact_kind: ManagedArtifactKind,
629 artifact_source: &Path,
630 ) -> Result<InstallReceipt, ModelManagementError> {
631 validate_model_id(model_id)?;
632 let shared_cache_references = match artifact_kind {
633 ManagedArtifactKind::Symlink => {
634 vec![artifact_source
635 .canonicalize()
636 .map_err(|source| ModelManagementError::Io {
637 path: artifact_source.to_path_buf(),
638 source,
639 })?]
640 }
641 ManagedArtifactKind::Directory => collect_shared_references(artifact_source, model_id)?,
642 ManagedArtifactKind::File => Vec::new(),
643 };
644 Ok(InstallReceipt {
645 model_id: model_id.to_string(),
646 managed_path,
647 artifact_kind,
648 source_model_id: source_model_id.to_string(),
649 source_revision: None,
650 creation_generation,
651 shared_cache_references,
652 adopted,
653 })
654 }
655
656 pub(crate) fn resume_install_intent(
657 &self,
658 model_id: &str,
659 ) -> Result<Option<InstallReceipt>, ModelManagementError> {
660 let Some(intent) = self.load_install_intent(model_id)? else {
661 return Ok(None);
662 };
663 if std::fs::symlink_metadata(&intent.receipt.managed_path).is_err() {
664 if let (Some(staging), Some(identity)) = (
665 intent.staging_path.as_deref(),
666 intent.staging_identity.as_ref(),
667 ) {
668 if std::fs::symlink_metadata(staging).is_ok() {
669 ensure_artifact_identity(staging, model_id, identity)?;
670 atomic_rename_noreplace(staging, &intent.receipt.managed_path).map_err(
671 |source| ModelManagementError::Io {
672 path: intent.receipt.managed_path.clone(),
673 source,
674 },
675 )?;
676 sync_directory(&self.models_dir).map_err(|source| {
677 ModelManagementError::Io {
678 path: self.models_dir.clone(),
679 source,
680 }
681 })?;
682 } else {
683 self.clear_install_intent(model_id)?;
684 return Ok(None);
685 }
686 } else {
687 self.clear_install_intent(model_id)?;
690 return Ok(None);
691 }
692 }
693 if let Some(identity) = intent.staging_identity.as_ref() {
694 ensure_artifact_identity(&intent.receipt.managed_path, model_id, identity)?;
695 }
696 self.validate_managed_artifact(&intent.receipt)?;
697 self.write_receipt(&intent.receipt)?;
698 self.clear_install_intent(model_id)?;
699 self.clear_tombstone(model_id)?;
700 Ok(Some(intent.receipt))
701 }
702
703 fn load_install_intent(
704 &self,
705 model_id: &str,
706 ) -> Result<Option<InstallJournal>, ModelManagementError> {
707 let path = self.install_journal_path(model_id);
708 let bytes = match read_private_state(&path) {
709 Ok(bytes) => bytes,
710 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
711 Err(source) => return Err(ModelManagementError::Io { path, source }),
712 };
713 let intent: InstallJournal =
714 serde_json::from_slice(&bytes).map_err(|error| ModelManagementError::InvalidState {
715 path: path.clone(),
716 message: error.to_string(),
717 })?;
718 ensure_identity(&path, model_id, &intent.model_id)?;
719 ensure_identity(&path, model_id, &intent.receipt.model_id)?;
720 validate_direct_child(&self.models_dir, &intent.receipt.managed_path, model_id)?;
721 if let Some(staging) = intent.staging_path.as_deref() {
722 validate_direct_child(&self.models_dir, staging, model_id)?;
723 if !staging
724 .file_name()
725 .and_then(|name| name.to_str())
726 .is_some_and(|name| name.starts_with(".car-install-"))
727 {
728 return Err(ModelManagementError::UnsafeManagedPath {
729 model_id: model_id.to_string(),
730 path: staging.to_path_buf(),
731 reason: "install journal staging path lacks CAR staging identity".into(),
732 });
733 }
734 }
735 Ok(Some(intent))
736 }
737
738 fn clear_install_intent(&self, model_id: &str) -> Result<(), ModelManagementError> {
739 let path = self.install_journal_path(model_id);
740 match std::fs::remove_file(&path) {
741 Ok(()) => sync_directory(&self.installs_root)
742 .map_err(|source| ModelManagementError::Io { path, source }),
743 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
744 Err(source) => Err(ModelManagementError::Io { path, source }),
745 }
746 }
747
748 pub(crate) fn materialize_managed_projection(
752 &self,
753 managed_leaf: &str,
754 source: &Path,
755 ) -> Result<PathBuf, ModelManagementError> {
756 validate_managed_leaf(managed_leaf, &self.models_dir)?;
757 create_private_dir(&self.models_dir)?;
758 let models_root = canonical_directory(&self.models_dir)?;
759 let source_parent =
760 source
761 .parent()
762 .ok_or_else(|| ModelManagementError::UnsafeManagedPath {
763 model_id: managed_leaf.to_string(),
764 path: source.to_path_buf(),
765 reason: "source has no parent".into(),
766 })?;
767 if canonical_directory(source_parent)? == models_root
768 && source == self.models_dir.join(managed_leaf)
769 {
770 return Err(ModelManagementError::UnsafeManagedPath {
771 model_id: managed_leaf.to_string(),
772 path: source.to_path_buf(),
773 reason: "pre-existing unreceipted managed projection cannot become CAR-owned"
774 .into(),
775 });
776 }
777 let target = source
778 .canonicalize()
779 .map_err(|source_error| ModelManagementError::Io {
780 path: source.to_path_buf(),
781 source: source_error,
782 })?;
783 let managed = self.models_dir.join(managed_leaf);
784 if std::fs::symlink_metadata(&managed).is_ok() {
785 return Err(ModelManagementError::UnsafeManagedPath {
786 model_id: managed_leaf.to_string(),
787 path: managed,
788 reason: "pre-existing managed projection has no CAR creation receipt".into(),
789 });
790 }
791 static NEXT_PROJECTION: AtomicU64 = AtomicU64::new(1);
792 let tmp = self.models_dir.join(format!(
793 ".managed-projection.{}.{}",
794 std::process::id(),
795 NEXT_PROJECTION.fetch_add(1, Ordering::Relaxed)
796 ));
797 create_symlink(&target, &tmp)?;
798 let result = atomic_rename_noreplace(&tmp, &managed).map_err(|source| {
799 ModelManagementError::UnsafeManagedPath {
800 model_id: managed_leaf.to_string(),
801 path: managed.clone(),
802 reason: format!("managed projection publication raced or failed: {source}"),
803 }
804 });
805 if result.is_err() {
806 let _ = std::fs::remove_file(&tmp);
807 }
808 result?;
809 sync_directory(&self.models_dir).map_err(|source| ModelManagementError::Io {
810 path: self.models_dir.clone(),
811 source,
812 })?;
813 Ok(managed)
814 }
815
816 pub(crate) fn create_install_staging(
817 &self,
818 model_id: &str,
819 ) -> Result<PathBuf, ModelManagementError> {
820 validate_model_id(model_id)?;
821 create_private_dir(&self.models_dir)?;
822 let digest = hex::encode(Sha256::digest(model_id.as_bytes()));
823 for _ in 0..16 {
824 let candidate = self.models_dir.join(format!(
825 ".car-install-{}-{:032x}",
826 &digest[..16],
827 rand::random::<u128>()
828 ));
829 match std::fs::create_dir(&candidate) {
830 Ok(()) => {
831 harden_private_directory(&candidate)?;
832 return Ok(candidate);
833 }
834 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
835 Err(source) => {
836 return Err(ModelManagementError::Io {
837 path: candidate,
838 source,
839 });
840 }
841 }
842 }
843 Err(ModelManagementError::InvalidState {
844 path: self.models_dir.clone(),
845 message: "could not allocate a collision-resistant install staging directory".into(),
846 })
847 }
848
849 pub(crate) fn publish_install_staging(
850 &self,
851 model_id: &str,
852 staging: &Path,
853 managed_leaf: &str,
854 ) -> Result<PathBuf, ModelManagementError> {
855 validate_model_id(model_id)?;
856 validate_direct_child(&self.models_dir, staging, model_id)?;
857 if !staging
858 .file_name()
859 .and_then(|name| name.to_str())
860 .is_some_and(|name| name.starts_with(".car-install-"))
861 {
862 return Err(ModelManagementError::UnsafeManagedPath {
863 model_id: model_id.to_string(),
864 path: staging.to_path_buf(),
865 reason: "install staging path lacks CAR staging identity".into(),
866 });
867 }
868 validate_managed_leaf(managed_leaf, &self.models_dir)?;
869 let managed = self.models_dir.join(managed_leaf);
870 atomic_rename_noreplace(staging, &managed).map_err(|source| {
871 ModelManagementError::UnsafeManagedPath {
872 model_id: model_id.to_string(),
873 path: managed.clone(),
874 reason: format!("could not atomically publish CAR-owned install: {source}"),
875 }
876 })?;
877 sync_directory(&self.models_dir).map_err(|source| ModelManagementError::Io {
878 path: self.models_dir.clone(),
879 source,
880 })?;
881 Ok(managed)
882 }
883
884 pub(crate) fn discard_install_staging(&self, staging: &Path) {
885 let _ = staging;
890 }
891
892 pub(crate) fn materialize_adopted_projection(
896 &self,
897 model_id: &str,
898 source: &Path,
899 ) -> Result<PathBuf, ModelManagementError> {
900 validate_model_id(model_id)?;
901 let managed = self.adopted_projection_path(model_id)?;
902 let leaf = managed
903 .file_name()
904 .and_then(|name| name.to_str())
905 .expect("adopted projection is an UTF-8 direct child");
906 self.materialize_managed_projection(leaf, source)
907 }
908
909 pub(crate) fn adopted_projection_path(
910 &self,
911 model_id: &str,
912 ) -> Result<PathBuf, ModelManagementError> {
913 validate_model_id(model_id)?;
914 let digest = hex::encode(Sha256::digest(model_id.as_bytes()));
915 Ok(self.models_dir.join(format!(".car-adopted-{digest}")))
916 }
917
918 fn remove_with_mutation(
919 &self,
920 model_id: &str,
921 removal_generation: u64,
922 ) -> Result<RemoveFromCarResult, ModelManagementError> {
923 let _activity = self.lock_idle_activity(model_id)?;
924
925 if let Some(journal) = self.load_removal_journal(model_id)? {
926 write_private_json(
927 &self.tombstone_path(model_id),
928 &RemovalTombstone {
929 model_id: model_id.to_string(),
930 removal_generation: journal.removal_generation,
931 artifact_kind: Some(journal.receipt.artifact_kind),
932 preserved_shared_cache_references: journal
933 .receipt
934 .shared_cache_references
935 .clone(),
936 },
937 )?;
938 return self.resume_removal_journal(journal);
939 }
940
941 let Some(receipt) = self.load_receipt(model_id)? else {
942 if let Some(tombstone) = self.load_tombstone(model_id)? {
943 if let Some(artifact_kind) = tombstone.artifact_kind {
944 return Ok(RemoveFromCarResult {
945 model_id: model_id.to_string(),
946 artifact_kind,
947 preserved_shared_cache_references: tombstone
948 .preserved_shared_cache_references,
949 });
950 }
951 }
952 return Err(ModelManagementError::MissingReceipt {
953 model_id: model_id.to_string(),
954 });
955 };
956 self.validate_managed_artifact(&receipt)?;
957 if receipt.artifact_kind == ManagedArtifactKind::Directory && !directory_removal_supported()
958 {
959 return Err(ModelManagementError::UnsafeManagedPath {
960 model_id: model_id.to_string(),
961 path: receipt.managed_path,
962 reason: "recursive directory removal is unsupported without object-bound traversal"
963 .into(),
964 });
965 }
966 let artifact_identity = artifact_identity(&receipt.managed_path, model_id)?;
967 let quarantine = self.allocate_quarantine(model_id)?;
968
969 write_private_json(
970 &self.tombstone_path(model_id),
971 &RemovalTombstone {
972 model_id: model_id.to_string(),
973 removal_generation,
974 artifact_kind: Some(receipt.artifact_kind),
975 preserved_shared_cache_references: receipt.shared_cache_references.clone(),
976 },
977 )?;
978 let journal = RemovalJournal {
979 model_id: model_id.to_string(),
980 removal_generation,
981 receipt,
982 quarantine_path: quarantine,
983 artifact_identity,
984 };
985 write_private_json(&self.removal_journal_path(model_id), &journal)?;
986 self.resume_removal_journal(journal)
987 }
988
989 fn lock_idle_activity(&self, model_id: &str) -> Result<OwnedFileLock, ModelManagementError> {
990 let mut activity = self.open_identity_lock(&self.lease_path(model_id), model_id)?;
991 validate_lock_identity(&mut activity, &self.lease_path(model_id), model_id)?;
992 match activity.try_lock() {
993 Ok(()) => Ok(OwnedFileLock::new(activity)),
994 Err(std::fs::TryLockError::WouldBlock) => Err(ModelManagementError::ModelInUse {
995 model_id: model_id.to_string(),
996 }),
997 Err(std::fs::TryLockError::Error(source)) => Err(ModelManagementError::Io {
998 path: self.lease_path(model_id),
999 source,
1000 }),
1001 }
1002 }
1003
1004 fn load_removal_journal(
1005 &self,
1006 model_id: &str,
1007 ) -> Result<Option<RemovalJournal>, ModelManagementError> {
1008 let path = self.removal_journal_path(model_id);
1009 let bytes = match read_private_state(&path) {
1010 Ok(bytes) => bytes,
1011 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
1012 Err(source) => return Err(ModelManagementError::Io { path, source }),
1013 };
1014 let journal: RemovalJournal =
1015 serde_json::from_slice(&bytes).map_err(|error| ModelManagementError::InvalidState {
1016 path: path.clone(),
1017 message: error.to_string(),
1018 })?;
1019 ensure_identity(&path, model_id, &journal.model_id)?;
1020 ensure_identity(&path, model_id, &journal.receipt.model_id)?;
1021 validate_direct_child(&self.models_dir, &journal.receipt.managed_path, model_id)?;
1022 if journal.receipt.artifact_kind == ManagedArtifactKind::Directory
1023 && !directory_removal_supported()
1024 {
1025 return Err(ModelManagementError::UnsafeManagedPath {
1026 model_id: model_id.to_string(),
1027 path: journal.receipt.managed_path,
1028 reason: "recursive directory removal journal is unsupported".into(),
1029 });
1030 }
1031 validate_direct_child(&self.models_dir, &journal.quarantine_path, model_id)?;
1032 if !journal
1033 .quarantine_path
1034 .file_name()
1035 .and_then(|name| name.to_str())
1036 .is_some_and(|name| name.starts_with(".car-remove-"))
1037 {
1038 return Err(ModelManagementError::UnsafeManagedPath {
1039 model_id: model_id.to_string(),
1040 path: journal.quarantine_path,
1041 reason: "removal journal has an invalid quarantine path".into(),
1042 });
1043 }
1044 Ok(Some(journal))
1045 }
1046
1047 fn allocate_quarantine(&self, model_id: &str) -> Result<PathBuf, ModelManagementError> {
1048 let digest = hex::encode(Sha256::digest(model_id.as_bytes()));
1049 for _ in 0..16 {
1050 let candidate = self.models_dir.join(format!(
1051 ".car-remove-{}-{:032x}",
1052 &digest[..16],
1053 rand::random::<u128>()
1054 ));
1055 if std::fs::symlink_metadata(&candidate)
1056 .is_err_and(|error| error.kind() == std::io::ErrorKind::NotFound)
1057 {
1058 return Ok(candidate);
1059 }
1060 }
1061 Err(ModelManagementError::InvalidState {
1062 path: self.models_dir.clone(),
1063 message: "could not allocate a collision-resistant removal quarantine".into(),
1064 })
1065 }
1066
1067 fn resume_removal_journal(
1068 &self,
1069 journal: RemovalJournal,
1070 ) -> Result<RemoveFromCarResult, ModelManagementError> {
1071 let model_id = &journal.model_id;
1072 if journal.receipt.artifact_kind == ManagedArtifactKind::Directory
1073 && !directory_removal_supported()
1074 {
1075 return Err(ModelManagementError::UnsafeManagedPath {
1076 model_id: model_id.clone(),
1077 path: journal.receipt.managed_path,
1078 reason: "recursive directory removal is unsupported without object-bound traversal"
1079 .into(),
1080 });
1081 }
1082 let original_exists = std::fs::symlink_metadata(&journal.receipt.managed_path).is_ok();
1083 let quarantine_exists = std::fs::symlink_metadata(&journal.quarantine_path).is_ok();
1084 let mut renamed_now = false;
1085 if original_exists && quarantine_exists {
1086 return Err(ModelManagementError::UnsafeManagedPath {
1087 model_id: model_id.clone(),
1088 path: journal.receipt.managed_path.clone(),
1089 reason: "both managed artifact and removal quarantine exist".into(),
1090 });
1091 }
1092 if original_exists {
1093 self.validate_managed_artifact(&journal.receipt)?;
1094 ensure_artifact_identity(
1095 &journal.receipt.managed_path,
1096 model_id,
1097 &journal.artifact_identity,
1098 )?;
1099 validate_direct_child(&self.models_dir, &journal.receipt.managed_path, model_id)?;
1102 ensure_artifact_identity(
1103 &journal.receipt.managed_path,
1104 model_id,
1105 &journal.artifact_identity,
1106 )?;
1107 atomic_rename_noreplace(&journal.receipt.managed_path, &journal.quarantine_path)
1108 .map_err(|source| ModelManagementError::Io {
1109 path: journal.receipt.managed_path.clone(),
1110 source,
1111 })?;
1112 sync_directory(&self.models_dir).map_err(|source| ModelManagementError::Io {
1113 path: self.models_dir.clone(),
1114 source,
1115 })?;
1116 renamed_now = true;
1117 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
1118 if let Some(hook) = &self.removal_hook {
1119 (hook.0)(RemovalPhase::AfterRename, &journal.quarantine_path);
1120 }
1121 }
1122
1123 if std::fs::symlink_metadata(&journal.quarantine_path).is_err() && renamed_now {
1124 return Err(ModelManagementError::UnsafeManagedPath {
1129 model_id: model_id.clone(),
1130 path: journal.quarantine_path,
1131 reason: QUARANTINE_VANISHED_REASON.into(),
1132 });
1133 }
1134
1135 if std::fs::symlink_metadata(&journal.quarantine_path).is_ok() {
1136 ensure_removal_identity(
1139 &journal.quarantine_path,
1140 model_id,
1141 &journal.artifact_identity,
1142 )?;
1143 match journal.receipt.artifact_kind {
1144 #[cfg(any(target_os = "macos", target_os = "linux"))]
1145 ManagedArtifactKind::Directory => {
1146 self.remove_quarantined_directory(
1147 &journal.quarantine_path,
1148 &journal.artifact_identity,
1149 model_id,
1150 )?;
1151 }
1152 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
1153 ManagedArtifactKind::Directory => {
1154 return Err(ModelManagementError::UnsafeManagedPath {
1155 model_id: model_id.clone(),
1156 path: journal.quarantine_path,
1157 reason:
1158 "recursive directory removal is unsupported without object-bound traversal"
1159 .into(),
1160 });
1161 }
1162 ManagedArtifactKind::Symlink | ManagedArtifactKind::File => {
1163 std::fs::remove_file(&journal.quarantine_path).map_err(|source| {
1164 ModelManagementError::Io {
1165 path: journal.quarantine_path.clone(),
1166 source,
1167 }
1168 })?;
1169 }
1170 }
1171 sync_directory(&self.models_dir).map_err(|source| ModelManagementError::Io {
1172 path: self.models_dir.clone(),
1173 source,
1174 })?;
1175 }
1176
1177 let receipt_path = self.receipt_path(model_id);
1178 match std::fs::remove_file(&receipt_path) {
1179 Ok(()) => {
1180 sync_directory(&self.receipts_root).map_err(|source| ModelManagementError::Io {
1181 path: self.receipts_root.clone(),
1182 source,
1183 })?
1184 }
1185 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1186 Err(source) => {
1187 return Err(ModelManagementError::Io {
1188 path: receipt_path,
1189 source,
1190 });
1191 }
1192 }
1193 let journal_path = self.removal_journal_path(model_id);
1194 match std::fs::remove_file(&journal_path) {
1195 Ok(()) => {
1196 sync_directory(&self.removals_root).map_err(|source| ModelManagementError::Io {
1197 path: self.removals_root.clone(),
1198 source,
1199 })?
1200 }
1201 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1202 Err(source) => {
1203 return Err(ModelManagementError::Io {
1204 path: journal_path,
1205 source,
1206 });
1207 }
1208 }
1209 Ok(RemoveFromCarResult {
1210 model_id: model_id.clone(),
1211 artifact_kind: journal.receipt.artifact_kind,
1212 preserved_shared_cache_references: journal.receipt.shared_cache_references,
1213 })
1214 }
1215
1216 fn resume_pending_quarantines(&self) {
1217 let Ok(entries) = std::fs::read_dir(&self.removals_root) else {
1218 return;
1219 };
1220 for entry in entries.flatten() {
1221 if !is_hashed_state_filename(&entry.file_name()) {
1222 continue;
1223 }
1224 let Ok(bytes) = read_private_state(&entry.path()) else {
1225 continue;
1226 };
1227 let Ok(candidate) = serde_json::from_slice::<RemovalJournal>(&bytes) else {
1228 continue;
1229 };
1230 if entry.path() != self.removal_journal_path(&candidate.model_id) {
1231 continue;
1232 }
1233 let Ok(guard) = self.begin_mutation(&candidate.model_id) else {
1234 continue;
1235 };
1236 let Ok(Some(journal)) = self.load_removal_journal(&candidate.model_id) else {
1237 continue;
1238 };
1239 if std::fs::symlink_metadata(&journal.receipt.managed_path).is_ok() {
1240 continue;
1241 }
1242 let _ = guard.resume_existing(journal);
1243 }
1244 }
1245
1246 fn open_identity_lock(
1247 &self,
1248 path: &Path,
1249 model_id: &str,
1250 ) -> Result<File, ModelManagementError> {
1251 let parent = path.parent().expect("lock path has parent");
1252 create_private_dir(parent)?;
1253 let initialization_path = parent.join(".initialization.lock");
1258 let initialization = match open_private_new(&initialization_path) {
1259 Ok(file) => file,
1260 Err(ModelManagementError::Io { source, .. })
1261 if source.kind() == std::io::ErrorKind::AlreadyExists =>
1262 {
1263 open_existing_identity_lock(&initialization_path)?
1264 }
1265 Err(error) => return Err(error),
1266 };
1267 initialization
1268 .lock()
1269 .map_err(|source| ModelManagementError::Io {
1270 path: initialization_path,
1271 source,
1272 })?;
1273 let _initialization = OwnedFileLock::new(initialization);
1274 match open_private_new(path) {
1277 Ok(file) => {
1278 #[cfg(test)]
1279 if let Some(hook) = self.identity_init_hook {
1280 hook(path);
1281 }
1282 file.lock().map_err(|source| ModelManagementError::Io {
1283 path: path.to_path_buf(),
1284 source,
1285 })?;
1286 let mut file = OwnedFileLock::new(file);
1287 let bytes = serde_json::to_vec(&IdentityRecord {
1288 model_id: model_id.to_string(),
1289 })
1290 .map_err(|error| ModelManagementError::InvalidState {
1291 path: path.to_path_buf(),
1292 message: error.to_string(),
1293 })?;
1294 file.write_all(&bytes)
1295 .and_then(|_| file.sync_all())
1296 .map_err(|source| ModelManagementError::Io {
1297 path: path.to_path_buf(),
1298 source,
1299 })?;
1300 return file
1301 .into_unlocked_file()
1302 .map_err(|source| ModelManagementError::Io {
1303 path: path.to_path_buf(),
1304 source,
1305 });
1306 }
1307 Err(ModelManagementError::Io { source, .. })
1308 if source.kind() == std::io::ErrorKind::AlreadyExists => {}
1309 Err(error) => return Err(error),
1310 }
1311 open_existing_identity_lock(path)
1312 }
1313
1314 fn validate_managed_artifact(
1315 &self,
1316 receipt: &InstallReceipt,
1317 ) -> Result<(), ModelManagementError> {
1318 let models_root = validate_models_root(&self.models_dir, &receipt.model_id)?;
1319 for hf_root in [crate::hf_cache::hf_home(), crate::hf_cache::hub_dir()] {
1322 let hf_root = crate::resource_policy::normalized_state_root_key(&hf_root);
1323 if models_root.starts_with(&hf_root) || hf_root.starts_with(&models_root) {
1324 return Err(ModelManagementError::UnsafeManagedPath {
1325 model_id: receipt.model_id.clone(),
1326 path: receipt.managed_path.clone(),
1327 reason: "managed models root overlaps the shared Hugging Face cache".into(),
1328 });
1329 }
1330 }
1331 let parent = receipt.managed_path.parent().ok_or_else(|| {
1332 ModelManagementError::UnsafeManagedPath {
1333 model_id: receipt.model_id.clone(),
1334 path: receipt.managed_path.clone(),
1335 reason: "managed path has no parent".into(),
1336 }
1337 })?;
1338 let canonical_parent = canonical_directory(parent)?;
1339 if canonical_parent != models_root || receipt.managed_path == models_root {
1340 return Err(ModelManagementError::UnsafeManagedPath {
1341 model_id: receipt.model_id.clone(),
1342 path: receipt.managed_path.clone(),
1343 reason: "artifact is not a direct child of the managed models root".into(),
1344 });
1345 }
1346 let actual = artifact_kind(&receipt.managed_path, &receipt.model_id)?;
1347 if actual != receipt.artifact_kind {
1348 return Err(ModelManagementError::UnsafeManagedPath {
1349 model_id: receipt.model_id.clone(),
1350 path: receipt.managed_path.clone(),
1351 reason: format!(
1352 "receipt says {:?}, filesystem is {:?}",
1353 receipt.artifact_kind, actual
1354 ),
1355 });
1356 }
1357 if actual == ManagedArtifactKind::Symlink {
1358 let target =
1359 receipt
1360 .managed_path
1361 .canonicalize()
1362 .map_err(|source| ModelManagementError::Io {
1363 path: receipt.managed_path.clone(),
1364 source,
1365 })?;
1366 if !receipt
1367 .shared_cache_references
1368 .iter()
1369 .any(|path| path == &target)
1370 {
1371 return Err(ModelManagementError::UnsafeManagedPath {
1372 model_id: receipt.model_id.clone(),
1373 path: receipt.managed_path.clone(),
1374 reason: "symlink target does not match a preserved shared-cache reference"
1375 .into(),
1376 });
1377 }
1378 } else {
1379 let current_references =
1380 collect_shared_references(&receipt.managed_path, &receipt.model_id)?;
1381 if current_references != receipt.shared_cache_references {
1382 return Err(ModelManagementError::UnsafeManagedPath {
1383 model_id: receipt.model_id.clone(),
1384 path: receipt.managed_path.clone(),
1385 reason: "shared-cache references no longer match the install receipt".into(),
1386 });
1387 }
1388 }
1389 Ok(())
1390 }
1391}
1392
1393struct OwnedFileLock(Option<File>);
1397
1398impl OwnedFileLock {
1399 fn new(file: File) -> Self {
1400 Self(Some(file))
1401 }
1402
1403 fn into_unlocked_file(mut self) -> std::io::Result<File> {
1404 self.unlock()?;
1405 Ok(self.0.take().expect("owned lock has a file"))
1406 }
1407}
1408
1409impl std::ops::Deref for OwnedFileLock {
1410 type Target = File;
1411
1412 fn deref(&self) -> &File {
1413 self.0.as_ref().expect("owned lock has a file")
1414 }
1415}
1416
1417impl std::ops::DerefMut for OwnedFileLock {
1418 fn deref_mut(&mut self) -> &mut File {
1419 self.0.as_mut().expect("owned lock has a file")
1420 }
1421}
1422
1423impl Drop for OwnedFileLock {
1424 fn drop(&mut self) {
1425 if let Some(file) = &self.0 {
1426 let _ = file.unlock();
1427 }
1428 }
1429}
1430
1431pub(crate) struct ModelMutationGuard {
1432 store: ModelManagementStore,
1433 model_id: String,
1434 _file: OwnedFileLock,
1435}
1436
1437impl ModelMutationGuard {
1438 pub(crate) fn remove(
1439 self,
1440 removal_generation: u64,
1441 ) -> Result<RemoveFromCarResult, ModelManagementError> {
1442 self.store
1443 .remove_with_mutation(&self.model_id, removal_generation)
1444 }
1445
1446 fn resume_existing(
1447 self,
1448 journal: RemovalJournal,
1449 ) -> Result<RemoveFromCarResult, ModelManagementError> {
1450 if journal.model_id != self.model_id {
1451 return Err(ModelManagementError::IdentityMismatch {
1452 path: self.store.removal_journal_path(&self.model_id),
1453 expected: self.model_id,
1454 actual: journal.model_id,
1455 });
1456 }
1457 let _activity = self.store.lock_idle_activity(&self.model_id)?;
1458 write_private_json(
1459 &self.store.tombstone_path(&self.model_id),
1460 &RemovalTombstone {
1461 model_id: self.model_id.clone(),
1462 removal_generation: journal.removal_generation,
1463 artifact_kind: Some(journal.receipt.artifact_kind),
1464 preserved_shared_cache_references: journal.receipt.shared_cache_references.clone(),
1465 },
1466 )?;
1467 self.store.resume_removal_journal(journal)
1468 }
1469}
1470
1471pub struct ModelLease {
1472 _file: OwnedFileLock,
1473}
1474
1475#[derive(Debug, Clone, PartialEq, Eq)]
1476pub struct RemoveFromCarResult {
1477 pub model_id: String,
1478 pub artifact_kind: ManagedArtifactKind,
1479 pub preserved_shared_cache_references: Vec<PathBuf>,
1480}
1481
1482#[derive(Debug, Serialize, Deserialize)]
1483#[serde(deny_unknown_fields)]
1484struct IdentityRecord {
1485 model_id: String,
1486}
1487
1488#[derive(Debug, Serialize, Deserialize)]
1489#[serde(deny_unknown_fields)]
1490struct RemovalTombstone {
1491 model_id: String,
1492 removal_generation: u64,
1493 #[serde(default)]
1494 artifact_kind: Option<ManagedArtifactKind>,
1495 #[serde(default)]
1496 preserved_shared_cache_references: Vec<PathBuf>,
1497}
1498
1499#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1500#[serde(deny_unknown_fields)]
1501struct ArtifactIdentity {
1502 kind: ManagedArtifactKind,
1503 #[serde(default)]
1504 device: u64,
1505 #[serde(default)]
1506 inode: u64,
1507 #[serde(default)]
1508 file_len: u64,
1509}
1510
1511#[derive(Debug, Clone, Serialize, Deserialize)]
1512#[serde(deny_unknown_fields)]
1513struct RemovalJournal {
1514 model_id: String,
1515 removal_generation: u64,
1516 receipt: InstallReceipt,
1517 quarantine_path: PathBuf,
1518 artifact_identity: ArtifactIdentity,
1519}
1520
1521#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1522#[serde(deny_unknown_fields)]
1523struct InstallJournal {
1524 model_id: String,
1525 receipt: InstallReceipt,
1526 #[serde(default)]
1527 staging_path: Option<PathBuf>,
1528 #[serde(default)]
1529 staging_identity: Option<ArtifactIdentity>,
1530}
1531
1532fn hashed_filename(model_id: &str) -> String {
1533 let digest = Sha256::digest(model_id.as_bytes());
1534 format!("{}.json", hex::encode(digest))
1535}
1536
1537fn is_hashed_state_filename(name: &std::ffi::OsStr) -> bool {
1538 let Some(name) = name.to_str() else {
1539 return false;
1540 };
1541 let Some(hex) = name.strip_suffix(".json") else {
1542 return false;
1543 };
1544 hex.len() == 64
1545 && hex
1546 .bytes()
1547 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1548}
1549
1550fn validate_model_id(model_id: &str) -> Result<(), ModelManagementError> {
1551 if !model_id.is_empty() && model_id.len() <= MAX_MODEL_ID_BYTES {
1552 return Ok(());
1553 }
1554 Err(ModelManagementError::InvalidState {
1555 path: PathBuf::new(),
1556 message: "model id must contain 1..=4096 UTF-8 bytes".into(),
1557 })
1558}
1559
1560fn validate_managed_leaf(
1561 managed_leaf: &str,
1562 models_dir: &Path,
1563) -> Result<(), ModelManagementError> {
1564 if Path::new(managed_leaf).components().count() == 1
1565 && matches!(
1566 Path::new(managed_leaf).components().next(),
1567 Some(std::path::Component::Normal(_))
1568 )
1569 {
1570 return Ok(());
1571 }
1572 Err(ModelManagementError::UnsafeManagedPath {
1573 model_id: managed_leaf.to_string(),
1574 path: models_dir.join(managed_leaf),
1575 reason: "managed model name must be one ordinary path component".into(),
1576 })
1577}
1578
1579fn validate_direct_child(
1580 root: &Path,
1581 candidate: &Path,
1582 model_id: &str,
1583) -> Result<(), ModelManagementError> {
1584 let canonical_root = validate_models_root(root, model_id)?;
1585 let parent = candidate
1586 .parent()
1587 .ok_or_else(|| ModelManagementError::UnsafeManagedPath {
1588 model_id: model_id.to_string(),
1589 path: candidate.to_path_buf(),
1590 reason: "managed path has no parent".into(),
1591 })?;
1592 if canonical_directory(parent)? != canonical_root {
1593 return Err(ModelManagementError::UnsafeManagedPath {
1594 model_id: model_id.to_string(),
1595 path: candidate.to_path_buf(),
1596 reason: "managed path is not a direct child of the models root".into(),
1597 });
1598 }
1599 Ok(())
1600}
1601
1602fn validate_models_root(root: &Path, model_id: &str) -> Result<PathBuf, ModelManagementError> {
1603 let metadata = std::fs::symlink_metadata(root).map_err(|source| ModelManagementError::Io {
1604 path: root.to_path_buf(),
1605 source,
1606 })?;
1607 if metadata.file_type().is_symlink() {
1608 return Err(ModelManagementError::UnsafeManagedPath {
1609 model_id: model_id.to_string(),
1610 path: root.to_path_buf(),
1611 reason: "managed models root must not be a symlink".into(),
1612 });
1613 }
1614 #[cfg(windows)]
1615 {
1616 use std::os::windows::fs::MetadataExt;
1617 const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x400;
1618 if metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0 {
1619 return Err(ModelManagementError::UnsafeManagedPath {
1620 model_id: model_id.to_string(),
1621 path: root.to_path_buf(),
1622 reason: "managed models root must not be a reparse point".into(),
1623 });
1624 }
1625 }
1626 if !metadata.is_dir() {
1627 return Err(ModelManagementError::UnsafeManagedPath {
1628 model_id: model_id.to_string(),
1629 path: root.to_path_buf(),
1630 reason: "managed models root is not a directory".into(),
1631 });
1632 }
1633 canonical_directory(root)
1634}
1635
1636fn ensure_identity(path: &Path, expected: &str, actual: &str) -> Result<(), ModelManagementError> {
1637 if expected == actual {
1638 return Ok(());
1639 }
1640 Err(ModelManagementError::IdentityMismatch {
1641 path: path.to_path_buf(),
1642 expected: expected.to_string(),
1643 actual: actual.to_string(),
1644 })
1645}
1646
1647fn validate_lock_identity(
1648 file: &mut File,
1649 path: &Path,
1650 model_id: &str,
1651) -> Result<(), ModelManagementError> {
1652 file.rewind().map_err(|source| ModelManagementError::Io {
1653 path: path.to_path_buf(),
1654 source,
1655 })?;
1656 let mut bytes = Vec::new();
1657 file.take(16 * 1024)
1658 .read_to_end(&mut bytes)
1659 .map_err(|source| ModelManagementError::Io {
1660 path: path.to_path_buf(),
1661 source,
1662 })?;
1663 let record: IdentityRecord =
1664 serde_json::from_slice(&bytes).map_err(|error| ModelManagementError::InvalidState {
1665 path: path.to_path_buf(),
1666 message: error.to_string(),
1667 })?;
1668 ensure_identity(path, model_id, &record.model_id)
1669}
1670
1671fn canonical_directory(path: &Path) -> Result<PathBuf, ModelManagementError> {
1672 path.canonicalize()
1673 .map_err(|source| ModelManagementError::Io {
1674 path: path.to_path_buf(),
1675 source,
1676 })
1677}
1678
1679fn artifact_identity(
1680 path: &Path,
1681 model_id: &str,
1682) -> Result<ArtifactIdentity, ModelManagementError> {
1683 let metadata = std::fs::symlink_metadata(path).map_err(|source| ModelManagementError::Io {
1684 path: path.to_path_buf(),
1685 source,
1686 })?;
1687 reject_windows_reparse(&metadata, path, model_id)?;
1688 let kind = if metadata.file_type().is_symlink() {
1689 ManagedArtifactKind::Symlink
1690 } else if metadata.is_dir() {
1691 ManagedArtifactKind::Directory
1692 } else if metadata.is_file() {
1693 ManagedArtifactKind::File
1694 } else {
1695 return Err(ModelManagementError::UnsafeManagedPath {
1696 model_id: model_id.to_string(),
1697 path: path.to_path_buf(),
1698 reason: "unsupported artifact type".into(),
1699 });
1700 };
1701 #[cfg(unix)]
1702 {
1703 use std::os::unix::fs::MetadataExt;
1704 Ok(ArtifactIdentity {
1705 kind,
1706 device: metadata.dev(),
1707 inode: metadata.ino(),
1708 file_len: metadata.len(),
1709 })
1710 }
1711 #[cfg(not(unix))]
1712 {
1713 Ok(ArtifactIdentity {
1714 kind,
1715 device: 0,
1716 inode: 0,
1717 file_len: metadata.len(),
1718 })
1719 }
1720}
1721
1722fn ensure_artifact_identity(
1723 path: &Path,
1724 model_id: &str,
1725 expected: &ArtifactIdentity,
1726) -> Result<(), ModelManagementError> {
1727 let actual = artifact_identity(path, model_id)?;
1728 if &actual == expected {
1729 return Ok(());
1730 }
1731 Err(ModelManagementError::UnsafeManagedPath {
1732 model_id: model_id.to_string(),
1733 path: path.to_path_buf(),
1734 reason: "managed artifact identity changed during removal".into(),
1735 })
1736}
1737
1738impl ArtifactIdentity {
1739 fn matches_for_removal(&self, actual: &ArtifactIdentity) -> bool {
1744 self.kind == actual.kind
1745 && self.device == actual.device
1746 && self.inode == actual.inode
1747 && (self.kind == ManagedArtifactKind::Directory || self.file_len == actual.file_len)
1748 }
1749}
1750
1751fn ensure_removal_identity(
1752 path: &Path,
1753 model_id: &str,
1754 expected: &ArtifactIdentity,
1755) -> Result<(), ModelManagementError> {
1756 let actual = artifact_identity(path, model_id)?;
1757 if expected.matches_for_removal(&actual) {
1758 return Ok(());
1759 }
1760 Err(ModelManagementError::UnsafeManagedPath {
1761 model_id: model_id.to_string(),
1762 path: path.to_path_buf(),
1763 reason: "managed artifact identity changed during removal".into(),
1764 })
1765}
1766
1767fn artifact_kind(path: &Path, model_id: &str) -> Result<ManagedArtifactKind, ModelManagementError> {
1768 let metadata = std::fs::symlink_metadata(path).map_err(|source| ModelManagementError::Io {
1769 path: path.to_path_buf(),
1770 source,
1771 })?;
1772 reject_windows_reparse(&metadata, path, model_id)?;
1773 if metadata.file_type().is_symlink() {
1774 Ok(ManagedArtifactKind::Symlink)
1775 } else if metadata.is_dir() {
1776 Ok(ManagedArtifactKind::Directory)
1777 } else if metadata.is_file() {
1778 Ok(ManagedArtifactKind::File)
1779 } else {
1780 Err(ModelManagementError::UnsafeManagedPath {
1781 model_id: model_id.to_string(),
1782 path: path.to_path_buf(),
1783 reason: "unsupported artifact type".into(),
1784 })
1785 }
1786}
1787
1788fn collect_shared_references(
1789 path: &Path,
1790 model_id: &str,
1791) -> Result<Vec<PathBuf>, ModelManagementError> {
1792 let metadata = std::fs::symlink_metadata(path).map_err(|source| ModelManagementError::Io {
1793 path: path.to_path_buf(),
1794 source,
1795 })?;
1796 if metadata.file_type().is_symlink() {
1797 return path
1798 .canonicalize()
1799 .map(|target| vec![target])
1800 .map_err(|source| ModelManagementError::Io {
1801 path: path.to_path_buf(),
1802 source,
1803 });
1804 }
1805 if !metadata.is_dir() {
1806 return Ok(Vec::new());
1807 }
1808 let mut references = Vec::new();
1809 for entry in std::fs::read_dir(path).map_err(|source| ModelManagementError::Io {
1810 path: path.to_path_buf(),
1811 source,
1812 })? {
1813 let entry = entry.map_err(|source| ModelManagementError::Io {
1814 path: path.to_path_buf(),
1815 source,
1816 })?;
1817 let entry_path = entry.path();
1818 let entry_metadata =
1819 std::fs::symlink_metadata(&entry_path).map_err(|source| ModelManagementError::Io {
1820 path: entry_path.clone(),
1821 source,
1822 })?;
1823 if entry_metadata.file_type().is_symlink() {
1824 let target = entry_path
1825 .canonicalize()
1826 .map_err(|source| ModelManagementError::Io {
1827 path: entry_path.clone(),
1828 source,
1829 })?;
1830 if std::fs::metadata(&target)
1831 .map_err(|source| ModelManagementError::Io {
1832 path: target.clone(),
1833 source,
1834 })?
1835 .is_dir()
1836 {
1837 return Err(ModelManagementError::UnsafeManagedPath {
1838 model_id: model_id.to_string(),
1839 path: entry_path,
1840 reason: "managed directory contains a directory symlink".into(),
1841 });
1842 }
1843 references.push(target);
1844 } else if entry_metadata.is_dir() {
1845 references.extend(collect_shared_references(&entry_path, model_id)?);
1846 }
1847 }
1848 references.sort();
1849 references.dedup();
1850 Ok(references)
1851}
1852
1853#[cfg(unix)]
1854fn create_symlink(target: &Path, link: &Path) -> Result<(), ModelManagementError> {
1855 std::os::unix::fs::symlink(target, link).map_err(|source| ModelManagementError::Io {
1856 path: link.to_path_buf(),
1857 source,
1858 })
1859}
1860
1861#[cfg(windows)]
1862fn create_symlink(target: &Path, link: &Path) -> Result<(), ModelManagementError> {
1863 let result = if target.is_dir() {
1864 std::os::windows::fs::symlink_dir(target, link)
1865 } else {
1866 std::os::windows::fs::symlink_file(target, link)
1867 };
1868 result.map_err(|source| ModelManagementError::Io {
1869 path: link.to_path_buf(),
1870 source,
1871 })
1872}
1873
1874fn create_private_dir(path: &Path) -> Result<(), ModelManagementError> {
1875 std::fs::create_dir_all(path).map_err(|source| ModelManagementError::Io {
1876 path: path.to_path_buf(),
1877 source,
1878 })?;
1879 harden_private_directory(path)?;
1880 Ok(())
1881}
1882
1883fn harden_private_directory(path: &Path) -> Result<(), ModelManagementError> {
1884 let metadata = std::fs::symlink_metadata(path).map_err(|source| ModelManagementError::Io {
1885 path: path.to_path_buf(),
1886 source,
1887 })?;
1888 if metadata.file_type().is_symlink() {
1889 return Err(ModelManagementError::UnsafeManagedPath {
1890 model_id: "model-management-state".into(),
1891 path: path.to_path_buf(),
1892 reason: "private model-management directory must not be a symlink".into(),
1893 });
1894 }
1895 reject_windows_reparse(&metadata, path, "model-management-state")?;
1896 #[cfg(unix)]
1897 {
1898 use std::os::unix::fs::PermissionsExt;
1899 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700)).map_err(
1900 |source| ModelManagementError::Io {
1901 path: path.to_path_buf(),
1902 source,
1903 },
1904 )?;
1905 }
1906 Ok(())
1907}
1908
1909fn reject_windows_reparse(
1910 metadata: &std::fs::Metadata,
1911 path: &Path,
1912 model_id: &str,
1913) -> Result<(), ModelManagementError> {
1914 #[cfg(windows)]
1915 {
1916 use std::os::windows::fs::MetadataExt;
1917 const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x400;
1918 if metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0 {
1919 return Err(ModelManagementError::UnsafeManagedPath {
1920 model_id: model_id.to_string(),
1921 path: path.to_path_buf(),
1922 reason: "model-management path is a Windows reparse point".into(),
1923 });
1924 }
1925 }
1926 #[cfg(not(windows))]
1927 let _ = (metadata, path, model_id);
1928 Ok(())
1929}
1930
1931fn open_private_new(path: &Path) -> Result<File, ModelManagementError> {
1932 let mut options = OpenOptions::new();
1933 options.read(true).write(true).create_new(true);
1934 #[cfg(unix)]
1935 {
1936 use std::os::unix::fs::OpenOptionsExt;
1937 options.mode(0o600);
1938 }
1939 options
1940 .open(path)
1941 .map_err(|source| ModelManagementError::Io {
1942 path: path.to_path_buf(),
1943 source,
1944 })
1945}
1946
1947fn open_existing_identity_lock(path: &Path) -> Result<File, ModelManagementError> {
1948 let metadata = std::fs::symlink_metadata(path).map_err(|source| ModelManagementError::Io {
1949 path: path.to_path_buf(),
1950 source,
1951 })?;
1952 if metadata.file_type().is_symlink() {
1953 return Err(ModelManagementError::InvalidState {
1954 path: path.to_path_buf(),
1955 message: "model-management lock sentinel must not be a symlink".into(),
1956 });
1957 }
1958 reject_windows_reparse(&metadata, path, "model-management-lock")?;
1959 #[cfg(unix)]
1960 {
1961 use std::os::unix::fs::MetadataExt;
1962 if metadata.nlink() != 1 {
1963 return Err(ModelManagementError::InvalidState {
1964 path: path.to_path_buf(),
1965 message: "model-management lock sentinel must not be hard-linked".into(),
1966 });
1967 }
1968 }
1969 let mut options = OpenOptions::new();
1970 options.read(true).write(true);
1971 #[cfg(unix)]
1972 {
1973 use std::os::unix::fs::OpenOptionsExt;
1974 options.custom_flags(libc::O_NOFOLLOW);
1975 }
1976 #[cfg(windows)]
1977 {
1978 use std::os::windows::fs::OpenOptionsExt;
1979 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
1980 options.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT);
1981 }
1982 let file = options
1983 .open(path)
1984 .map_err(|source| ModelManagementError::Io {
1985 path: path.to_path_buf(),
1986 source,
1987 })?;
1988 let opened_metadata = file.metadata().map_err(|source| ModelManagementError::Io {
1989 path: path.to_path_buf(),
1990 source,
1991 })?;
1992 if !opened_metadata.is_file() {
1993 return Err(ModelManagementError::InvalidState {
1994 path: path.to_path_buf(),
1995 message: "model-management lock sentinel is not a regular file".into(),
1996 });
1997 }
1998 reject_windows_reparse(&opened_metadata, path, "model-management-lock")?;
1999 #[cfg(unix)]
2000 {
2001 use std::os::unix::fs::MetadataExt;
2002 if opened_metadata.nlink() != 1 {
2003 return Err(ModelManagementError::InvalidState {
2004 path: path.to_path_buf(),
2005 message: "opened model-management lock sentinel is hard-linked".into(),
2006 });
2007 }
2008 }
2009 Ok(file)
2010}
2011
2012pub fn read_receipt(
2016 state_root: &Path,
2017 model_id: &str,
2018) -> Result<Option<InstallReceipt>, ModelManagementError> {
2019 validate_model_id(model_id)?;
2020 let state_root = crate::resource_policy::normalized_state_root_key(state_root);
2021 read_receipt_at(
2022 &receipts_dir(&state_root).join(hashed_filename(model_id)),
2023 model_id,
2024 )
2025}
2026
2027fn receipts_dir(state_root: &Path) -> PathBuf {
2030 state_root.join("model-management").join("receipts")
2031}
2032
2033fn read_receipt_at(
2034 path: &Path,
2035 model_id: &str,
2036) -> Result<Option<InstallReceipt>, ModelManagementError> {
2037 let bytes = match read_private_state(path) {
2038 Ok(bytes) => bytes,
2039 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
2040 Err(source) => {
2041 return Err(ModelManagementError::Io {
2042 path: path.to_path_buf(),
2043 source,
2044 })
2045 }
2046 };
2047 let receipt: InstallReceipt =
2048 serde_json::from_slice(&bytes).map_err(|error| ModelManagementError::InvalidState {
2049 path: path.to_path_buf(),
2050 message: error.to_string(),
2051 })?;
2052 ensure_identity(path, model_id, &receipt.model_id)?;
2053 Ok(Some(receipt))
2054}
2055
2056fn read_private_state(path: &Path) -> std::io::Result<Vec<u8>> {
2057 let mut options = OpenOptions::new();
2058 options.read(true);
2059 #[cfg(unix)]
2060 {
2061 use std::os::unix::fs::OpenOptionsExt;
2062 options.custom_flags(libc::O_NOFOLLOW);
2063 }
2064 let file = options.open(path)?;
2065 let metadata = file.metadata()?;
2066 if !metadata.is_file() || metadata.len() > MAX_STATE_BYTES {
2067 return Err(std::io::Error::new(
2068 std::io::ErrorKind::InvalidData,
2069 "model-management state must be a bounded regular file",
2070 ));
2071 }
2072 #[cfg(unix)]
2073 {
2074 use std::os::unix::fs::MetadataExt;
2075 if metadata.nlink() != 1 {
2076 return Err(std::io::Error::new(
2077 std::io::ErrorKind::InvalidData,
2078 "model-management state must not be hard-linked",
2079 ));
2080 }
2081 }
2082 let mut bytes = Vec::with_capacity(metadata.len() as usize);
2083 file.take(MAX_STATE_BYTES + 1).read_to_end(&mut bytes)?;
2084 if bytes.len() as u64 > MAX_STATE_BYTES {
2085 return Err(std::io::Error::new(
2086 std::io::ErrorKind::InvalidData,
2087 "model-management state exceeds the size limit",
2088 ));
2089 }
2090 Ok(bytes)
2091}
2092
2093fn write_private_json<T: Serialize>(path: &Path, value: &T) -> Result<(), ModelManagementError> {
2094 let _guard = state_mutation_lock()
2095 .lock()
2096 .unwrap_or_else(std::sync::PoisonError::into_inner);
2097 let parent = path
2098 .parent()
2099 .expect("model-management record path has a parent");
2100 create_private_dir(parent)?;
2101 let tmp = (0..16)
2102 .map(|_| {
2103 parent.join(format!(
2104 ".{}.{}-{:032x}.tmp",
2105 std::process::id(),
2106 path.file_name()
2107 .and_then(|name| name.to_str())
2108 .unwrap_or("state"),
2109 rand::random::<u128>()
2110 ))
2111 })
2112 .find(|candidate| {
2113 std::fs::symlink_metadata(candidate)
2114 .is_err_and(|error| error.kind() == std::io::ErrorKind::NotFound)
2115 })
2116 .ok_or_else(|| ModelManagementError::InvalidState {
2117 path: parent.to_path_buf(),
2118 message: "could not allocate collision-resistant state staging file".into(),
2119 })?;
2120 let bytes =
2121 serde_json::to_vec_pretty(value).map_err(|error| ModelManagementError::InvalidState {
2122 path: path.to_path_buf(),
2123 message: error.to_string(),
2124 })?;
2125 let result = (|| {
2126 let mut file = open_private_new(&tmp)?;
2127 file.write_all(&bytes)
2128 .and_then(|_| file.sync_all())
2129 .map_err(|source| ModelManagementError::Io {
2130 path: tmp.clone(),
2131 source,
2132 })?;
2133 atomic_replace(&tmp, path).map_err(|source| ModelManagementError::Io {
2134 path: path.to_path_buf(),
2135 source,
2136 })?;
2137 car_secrets::harden_owner_only(path);
2138 sync_directory(parent).map_err(|source| ModelManagementError::Io {
2139 path: parent.to_path_buf(),
2140 source,
2141 })
2142 })();
2143 if result.is_err() {
2144 let _ = std::fs::remove_file(&tmp);
2145 }
2146 result
2147}
2148
2149fn state_mutation_lock() -> &'static Mutex<()> {
2150 static LOCK: OnceLock<Mutex<()>> = OnceLock::new();
2151 LOCK.get_or_init(|| Mutex::new(()))
2152}
2153
2154#[cfg(not(windows))]
2155fn atomic_replace(source: &Path, destination: &Path) -> std::io::Result<()> {
2156 std::fs::rename(source, destination)
2157}
2158
2159#[cfg(target_os = "macos")]
2160fn atomic_rename_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2161 use std::ffi::CString;
2162 use std::os::unix::ffi::OsStrExt;
2163
2164 let source = CString::new(source.as_os_str().as_bytes())
2165 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?;
2166 let destination = CString::new(destination.as_os_str().as_bytes())
2167 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?;
2168 let result =
2170 unsafe { libc::renamex_np(source.as_ptr(), destination.as_ptr(), libc::RENAME_EXCL) };
2171 if result == 0 {
2172 Ok(())
2173 } else {
2174 Err(std::io::Error::last_os_error())
2175 }
2176}
2177
2178#[cfg(target_os = "linux")]
2179fn atomic_rename_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2180 use std::ffi::CString;
2181 use std::os::unix::ffi::OsStrExt;
2182
2183 let source = CString::new(source.as_os_str().as_bytes())
2184 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?;
2185 let destination = CString::new(destination.as_os_str().as_bytes())
2186 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?;
2187 let result = unsafe {
2189 libc::renameat2(
2190 libc::AT_FDCWD,
2191 source.as_ptr(),
2192 libc::AT_FDCWD,
2193 destination.as_ptr(),
2194 libc::RENAME_NOREPLACE,
2195 )
2196 };
2197 if result == 0 {
2198 Ok(())
2199 } else {
2200 Err(std::io::Error::last_os_error())
2201 }
2202}
2203
2204#[cfg(windows)]
2205fn atomic_rename_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2206 use std::os::windows::ffi::OsStrExt;
2207
2208 const MOVEFILE_WRITE_THROUGH: u32 = 0x8;
2209 #[link(name = "kernel32")]
2210 unsafe extern "system" {
2211 fn MoveFileExW(existing: *const u16, new: *const u16, flags: u32) -> i32;
2212 }
2213 let source = source
2214 .as_os_str()
2215 .encode_wide()
2216 .chain(Some(0))
2217 .collect::<Vec<_>>();
2218 let destination = destination
2219 .as_os_str()
2220 .encode_wide()
2221 .chain(Some(0))
2222 .collect::<Vec<_>>();
2223 let moved = unsafe {
2225 MoveFileExW(
2226 source.as_ptr(),
2227 destination.as_ptr(),
2228 MOVEFILE_WRITE_THROUGH,
2229 )
2230 };
2231 if moved == 0 {
2232 Err(std::io::Error::last_os_error())
2233 } else {
2234 Ok(())
2235 }
2236}
2237
2238#[cfg(not(any(target_os = "macos", target_os = "linux", windows)))]
2239fn atomic_rename_noreplace(_source: &Path, _destination: &Path) -> std::io::Result<()> {
2240 Err(std::io::Error::new(
2241 std::io::ErrorKind::Unsupported,
2242 "atomic no-replace publication is unavailable on this platform",
2243 ))
2244}
2245
2246#[cfg(windows)]
2247fn atomic_replace(source: &Path, destination: &Path) -> std::io::Result<()> {
2248 use std::os::windows::ffi::OsStrExt;
2249
2250 const MOVEFILE_REPLACE_EXISTING: u32 = 0x1;
2251 const MOVEFILE_WRITE_THROUGH: u32 = 0x8;
2252 #[link(name = "kernel32")]
2253 unsafe extern "system" {
2254 fn MoveFileExW(existing: *const u16, new: *const u16, flags: u32) -> i32;
2255 }
2256 let source = source
2257 .as_os_str()
2258 .encode_wide()
2259 .chain(Some(0))
2260 .collect::<Vec<_>>();
2261 let destination = destination
2262 .as_os_str()
2263 .encode_wide()
2264 .chain(Some(0))
2265 .collect::<Vec<_>>();
2266 let replaced = unsafe {
2269 MoveFileExW(
2270 source.as_ptr(),
2271 destination.as_ptr(),
2272 MOVEFILE_REPLACE_EXISTING | MOVEFILE_WRITE_THROUGH,
2273 )
2274 };
2275 if replaced == 0 {
2276 Err(std::io::Error::last_os_error())
2277 } else {
2278 Ok(())
2279 }
2280}
2281
2282#[cfg(unix)]
2283fn sync_directory(path: &Path) -> std::io::Result<()> {
2284 File::open(path)?.sync_all()
2285}
2286
2287#[cfg(any(target_os = "macos", target_os = "linux"))]
2300mod object_bound {
2301 use std::ffi::{CStr, CString};
2302 use std::fs::File;
2303 use std::os::fd::{AsRawFd, FromRawFd};
2304 use std::os::unix::ffi::OsStrExt;
2305 use std::os::unix::fs::MetadataExt;
2306 use std::path::Path;
2307
2308 pub(super) const MAX_REMOVAL_DEPTH: usize = 32;
2309
2310 const DIRECTORY_FLAGS: libc::c_int =
2311 libc::O_RDONLY | libc::O_DIRECTORY | libc::O_CLOEXEC | libc::O_NOFOLLOW;
2312
2313 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
2314 pub(super) struct Identity {
2315 pub(super) device: u64,
2316 pub(super) inode: u64,
2317 }
2318
2319 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
2320 pub(super) enum EntryKind {
2321 Directory,
2322 RegularFile,
2323 Symlink,
2324 Other,
2325 }
2326
2327 pub(super) fn c_name(path: &Path) -> std::io::Result<CString> {
2328 CString::new(path.as_os_str().as_bytes())
2329 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))
2330 }
2331
2332 pub(super) fn open_directory(path: &Path) -> std::io::Result<File> {
2334 let name = c_name(path)?;
2335 let raw = unsafe { libc::open(name.as_ptr(), DIRECTORY_FLAGS) };
2337 if raw < 0 {
2338 return Err(std::io::Error::last_os_error());
2339 }
2340 Ok(unsafe { File::from_raw_fd(raw) })
2342 }
2343
2344 #[cfg(target_os = "macos")]
2348 pub(super) fn open_child_directory(parent: &File, name: &CStr) -> std::io::Result<File> {
2349 let raw = unsafe { libc::openat(parent.as_raw_fd(), name.as_ptr(), DIRECTORY_FLAGS) };
2351 if raw < 0 {
2352 return Err(std::io::Error::last_os_error());
2353 }
2354 Ok(unsafe { File::from_raw_fd(raw) })
2356 }
2357
2358 #[cfg(target_os = "linux")]
2362 pub(super) fn open_child_directory(parent: &File, name: &CStr) -> std::io::Result<File> {
2363 let mut how: libc::open_how = unsafe { std::mem::zeroed() };
2365 how.flags = DIRECTORY_FLAGS as u64;
2366 how.resolve = libc::RESOLVE_BENEATH | libc::RESOLVE_NO_SYMLINKS | libc::RESOLVE_NO_XDEV;
2367 let raw = unsafe {
2370 libc::syscall(
2371 libc::SYS_openat2,
2372 parent.as_raw_fd(),
2373 name.as_ptr(),
2374 &how as *const libc::open_how,
2375 std::mem::size_of::<libc::open_how>(),
2376 )
2377 };
2378 if raw < 0 {
2379 return Err(std::io::Error::last_os_error());
2380 }
2381 Ok(unsafe { File::from_raw_fd(raw as libc::c_int) })
2383 }
2384
2385 pub(super) fn same_filesystem(root: &Identity, child: &Identity) -> bool {
2389 root.device == child.device
2390 }
2391
2392 pub(super) fn is_not_a_directory(error: &std::io::Error) -> bool {
2393 matches!(
2394 error.raw_os_error(),
2395 Some(libc::ENOTDIR) | Some(libc::ELOOP)
2396 )
2397 }
2398
2399 pub(super) fn identity_of(file: &File) -> std::io::Result<Identity> {
2400 let metadata = file.metadata()?;
2401 Ok(Identity {
2402 device: metadata.dev(),
2403 inode: metadata.ino(),
2404 })
2405 }
2406
2407 #[allow(clippy::unnecessary_cast)]
2412 pub(super) fn stat_entry(parent: &File, name: &CStr) -> std::io::Result<(EntryKind, Identity)> {
2413 let mut stat: libc::stat = unsafe { std::mem::zeroed() };
2415 let result = unsafe {
2418 libc::fstatat(
2419 parent.as_raw_fd(),
2420 name.as_ptr(),
2421 &mut stat,
2422 libc::AT_SYMLINK_NOFOLLOW,
2423 )
2424 };
2425 if result != 0 {
2426 return Err(std::io::Error::last_os_error());
2427 }
2428 let kind = match stat.st_mode & libc::S_IFMT {
2429 libc::S_IFDIR => EntryKind::Directory,
2430 libc::S_IFREG => EntryKind::RegularFile,
2431 libc::S_IFLNK => EntryKind::Symlink,
2432 _ => EntryKind::Other,
2433 };
2434 Ok((
2435 kind,
2436 Identity {
2437 device: stat.st_dev as u64,
2438 inode: stat.st_ino as u64,
2439 },
2440 ))
2441 }
2442
2443 pub(super) fn list_names(directory: &File) -> std::io::Result<Vec<CString>> {
2446 let dot = c".";
2447 let raw = unsafe { libc::openat(directory.as_raw_fd(), dot.as_ptr(), DIRECTORY_FLAGS) };
2449 if raw < 0 {
2450 return Err(std::io::Error::last_os_error());
2451 }
2452 let stream = unsafe { libc::fdopendir(raw) };
2454 if stream.is_null() {
2455 let error = std::io::Error::last_os_error();
2456 unsafe {
2458 libc::close(raw);
2459 }
2460 return Err(error);
2461 }
2462 let mut names = Vec::new();
2463 let enumeration = loop {
2464 set_errno(0);
2465 let entry = unsafe { libc::readdir(stream) };
2467 if entry.is_null() {
2468 let errno = errno();
2469 break if errno == 0 {
2470 Ok(())
2471 } else {
2472 Err(std::io::Error::from_raw_os_error(errno))
2473 };
2474 }
2475 let bytes = unsafe { CStr::from_ptr((*entry).d_name.as_ptr()) }.to_bytes();
2477 if bytes != b"." && bytes != b".." {
2478 match CString::new(bytes) {
2479 Ok(name) => names.push(name),
2480 Err(error) => {
2481 break Err(std::io::Error::new(std::io::ErrorKind::InvalidData, error))
2482 }
2483 }
2484 }
2485 };
2486 let closed = unsafe { libc::closedir(stream) };
2488 enumeration?;
2489 if closed < 0 {
2490 return Err(std::io::Error::last_os_error());
2491 }
2492 names.sort();
2493 Ok(names)
2494 }
2495
2496 pub(super) fn unlink_entry(parent: &File, name: &CStr) -> std::io::Result<()> {
2498 let result = unsafe { libc::unlinkat(parent.as_raw_fd(), name.as_ptr(), 0) };
2500 if result == 0 {
2501 return Ok(());
2502 }
2503 let error = std::io::Error::last_os_error();
2504 if error.kind() == std::io::ErrorKind::NotFound {
2505 return Ok(());
2506 }
2507 Err(error)
2508 }
2509
2510 pub(super) fn remove_empty_directory(parent: &File, name: &CStr) -> std::io::Result<()> {
2513 let result =
2515 unsafe { libc::unlinkat(parent.as_raw_fd(), name.as_ptr(), libc::AT_REMOVEDIR) };
2516 if result == 0 {
2517 return Ok(());
2518 }
2519 let error = std::io::Error::last_os_error();
2520 if error.kind() == std::io::ErrorKind::NotFound {
2521 return Ok(());
2522 }
2523 Err(error)
2524 }
2525
2526 pub(super) fn is_not_empty(error: &std::io::Error) -> bool {
2527 matches!(
2528 error.raw_os_error(),
2529 Some(libc::ENOTEMPTY) | Some(libc::EEXIST)
2530 )
2531 }
2532
2533 #[cfg(target_os = "macos")]
2534 fn errno() -> i32 {
2535 unsafe { *libc::__error() }
2537 }
2538
2539 #[cfg(target_os = "macos")]
2540 fn set_errno(value: i32) {
2541 unsafe {
2543 *libc::__error() = value;
2544 }
2545 }
2546
2547 #[cfg(target_os = "linux")]
2548 fn errno() -> i32 {
2549 unsafe { *libc::__errno_location() }
2551 }
2552
2553 #[cfg(target_os = "linux")]
2554 fn set_errno(value: i32) {
2555 unsafe {
2557 *libc::__errno_location() = value;
2558 }
2559 }
2560}
2561
2562#[derive(Debug, Clone, Serialize, Deserialize)]
2566struct DetachedTreeJournal {
2567 source: PathBuf,
2568 quarantine: PathBuf,
2569 identity: ArtifactIdentity,
2570}
2571
2572#[derive(Debug)]
2575pub(crate) struct DetachError {
2576 pub(crate) detached: bool,
2577 pub(crate) error: ModelManagementError,
2578}
2579
2580#[derive(Debug, Serialize, Deserialize)]
2584struct RetiringIntent {
2585 model_id: String,
2586}
2587
2588pub(crate) const RETIRE_QUARANTINE_PREFIX: &str = ".car-retire-";
2590
2591impl ModelManagementStore {
2592 fn detached_tree_journals(&self) -> PathBuf {
2593 self.management_state_dir.join("hub-removals")
2594 }
2595
2596 fn retiring_intent_path(&self, model_id: &str) -> PathBuf {
2597 self.management_state_dir
2598 .join("retiring")
2599 .join(hashed_filename(model_id))
2600 }
2601
2602 pub(crate) fn note_retiring(&self, model_id: &str) -> Result<(), ModelManagementError> {
2604 validate_model_id(model_id)?;
2605 write_private_json(
2606 &self.retiring_intent_path(model_id),
2607 &RetiringIntent {
2608 model_id: model_id.to_string(),
2609 },
2610 )
2611 }
2612
2613 pub(crate) fn retired(&self, model_id: &str) {
2615 let _ = std::fs::remove_file(self.retiring_intent_path(model_id));
2616 }
2617
2618 pub fn resume_retiring(&self) -> Vec<ModelManagementError> {
2623 let mut errors = Vec::new();
2624 let Ok(read) = std::fs::read_dir(self.management_state_dir.join("retiring")) else {
2625 return errors;
2626 };
2627 for entry in read.filter_map(Result::ok) {
2628 let path = entry.path();
2629 let Some(intent) = read_private_state(&path)
2630 .ok()
2631 .and_then(|bytes| serde_json::from_slice::<RetiringIntent>(&bytes).ok())
2632 .filter(|intent| {
2633 path.file_name() == Some(hashed_filename(&intent.model_id).as_ref())
2634 })
2635 else {
2636 tracing::warn!(path = %path.display(), "unreadable retirement intent left in place");
2637 continue;
2638 };
2639 let id = &intent.model_id;
2640 let receipt_present = self.receipt_path(id).exists();
2641 if receipt_present {
2642 if !self.removal_journal_path(id).exists() {
2643 self.retired(id);
2645 }
2646 continue;
2647 }
2648 match self.clear_tombstone(id) {
2649 Ok(()) => self.retired(id),
2650 Err(error) => errors.push(error),
2651 }
2652 }
2653 errors
2654 }
2655
2656 pub(crate) fn detach_and_remove_tree(
2670 &self,
2671 source: &Path,
2672 quarantine: &Path,
2673 label: &str,
2674 verify: impl FnOnce() -> Result<(), String>,
2675 ) -> Result<(), DetachError> {
2676 let before = |error| DetachError {
2677 detached: false,
2678 error,
2679 };
2680 if !directory_removal_supported() {
2681 return Err(before(ModelManagementError::UnsafeManagedPath {
2682 model_id: label.to_string(),
2683 path: source.to_path_buf(),
2684 reason: "directory removal is not supported on this platform".into(),
2685 }));
2686 }
2687 let identity = artifact_identity(source, label).map_err(before)?;
2688 if identity.kind != ManagedArtifactKind::Directory {
2689 return Err(before(ModelManagementError::UnsafeManagedPath {
2690 model_id: label.to_string(),
2691 path: source.to_path_buf(),
2692 reason: "a retired repo must be a directory".into(),
2693 }));
2694 }
2695 if quarantine.parent() != source.parent() {
2696 return Err(before(ModelManagementError::UnsafeManagedPath {
2697 model_id: label.to_string(),
2698 path: quarantine.to_path_buf(),
2699 reason: "the quarantine must sit beside the directory it detaches".into(),
2700 }));
2701 }
2702 verify().map_err(|reason| {
2703 before(ModelManagementError::UnsafeManagedPath {
2704 model_id: label.to_string(),
2705 path: source.to_path_buf(),
2706 reason,
2707 })
2708 })?;
2709 let again = artifact_identity(source, label).map_err(before)?;
2710 if again.device != identity.device || again.inode != identity.inode {
2711 return Err(before(ModelManagementError::UnsafeManagedPath {
2712 model_id: label.to_string(),
2713 path: source.to_path_buf(),
2714 reason: "the repo directory changed while it was being checked".into(),
2715 }));
2716 }
2717 let journal = DetachedTreeJournal {
2718 source: source.to_path_buf(),
2719 quarantine: quarantine.to_path_buf(),
2720 identity,
2721 };
2722 let journal_path = self
2723 .detached_tree_journals()
2724 .join(hashed_filename(&quarantine.to_string_lossy()));
2725 if let Some(parent) = journal_path.parent() {
2726 std::fs::create_dir_all(parent)
2727 .map_err(|source| ModelManagementError::Io {
2728 path: parent.to_path_buf(),
2729 source,
2730 })
2731 .map_err(before)?;
2732 }
2733 write_private_json(&journal_path, &journal).map_err(before)?;
2734 if let Err(io) = atomic_rename_noreplace(source, quarantine) {
2735 let _ = std::fs::remove_file(&journal_path);
2736 return Err(before(ModelManagementError::Io {
2737 path: source.to_path_buf(),
2738 source: io,
2739 }));
2740 }
2741 self.finish_detached_tree(&journal, &journal_path, label)
2742 .map_err(|error| DetachError {
2743 detached: true,
2744 error,
2745 })
2746 }
2747
2748 #[cfg(any(target_os = "macos", target_os = "linux"))]
2749 fn finish_detached_tree(
2750 &self,
2751 journal: &DetachedTreeJournal,
2752 journal_path: &Path,
2753 label: &str,
2754 ) -> Result<(), ModelManagementError> {
2755 self.remove_quarantined_directory(&journal.quarantine, &journal.identity, label)?;
2756 let _ = std::fs::remove_file(journal_path);
2757 Ok(())
2758 }
2759
2760 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
2761 fn finish_detached_tree(
2762 &self,
2763 journal: &DetachedTreeJournal,
2764 _journal_path: &Path,
2765 label: &str,
2766 ) -> Result<(), ModelManagementError> {
2767 Err(ModelManagementError::UnsafeManagedPath {
2768 model_id: label.to_string(),
2769 path: journal.quarantine.clone(),
2770 reason: "directory removal is not supported on this platform".into(),
2771 })
2772 }
2773
2774 pub fn resume_detached_trees(&self) -> Vec<ModelManagementError> {
2778 let mut errors = Vec::new();
2779 let Ok(read) = std::fs::read_dir(self.detached_tree_journals()) else {
2780 return errors;
2781 };
2782 for entry in read.filter_map(Result::ok) {
2783 let path = entry.path();
2784 let journal: DetachedTreeJournal = match read_private_state(&path)
2785 .ok()
2786 .and_then(|bytes| serde_json::from_slice(&bytes).ok())
2787 {
2788 Some(journal) => journal,
2789 None => {
2790 tracing::warn!(path = %path.display(), "unreadable retirement journal left in place");
2791 continue;
2792 }
2793 };
2794 let named = journal
2799 .quarantine
2800 .file_name()
2801 .and_then(|name| name.to_str())
2802 .is_some_and(|name| name.starts_with(RETIRE_QUARANTINE_PREFIX));
2803 if !named || journal.quarantine.parent() != journal.source.parent() {
2804 tracing::warn!(path = %path.display(), "retirement journal names a foreign quarantine; left in place");
2805 continue;
2806 }
2807 match std::fs::symlink_metadata(&journal.quarantine) {
2808 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2810 let _ = std::fs::remove_file(&path);
2811 continue;
2812 }
2813 Err(source) => {
2814 errors.push(ModelManagementError::Io {
2815 path: journal.quarantine.clone(),
2816 source,
2817 });
2818 continue;
2819 }
2820 Ok(_) => {}
2821 }
2822 if let Err(error) = self.finish_detached_tree(&journal, &path, "retired repo") {
2823 errors.push(error);
2824 }
2825 }
2826 errors
2827 }
2828}
2829
2830#[cfg(any(target_os = "macos", target_os = "linux"))]
2831impl ModelManagementStore {
2832 fn remove_quarantined_directory(
2835 &self,
2836 quarantine: &Path,
2837 expected: &ArtifactIdentity,
2838 model_id: &str,
2839 ) -> Result<(), ModelManagementError> {
2840 use object_bound::{EntryKind, Identity, MAX_REMOVAL_DEPTH};
2841 use std::ffi::CString;
2842 use std::fs::File;
2843
2844 struct CapturedDirectory {
2845 file: File,
2846 name: CString,
2847 parent: usize,
2848 depth: usize,
2849 identity: Identity,
2850 }
2851 struct CapturedEntry {
2852 parent: usize,
2853 name: CString,
2854 identity: Identity,
2855 }
2856
2857 let unsafe_path = |reason: &str| ModelManagementError::UnsafeManagedPath {
2858 model_id: model_id.to_string(),
2859 path: quarantine.to_path_buf(),
2860 reason: reason.into(),
2861 };
2862 let io_error = |source: std::io::Error| ModelManagementError::Io {
2863 path: quarantine.to_path_buf(),
2864 source,
2865 };
2866
2867 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
2868 if let Some(hook) = &self.removal_hook {
2869 (hook.0)(RemovalPhase::BeforeRootOpen, quarantine);
2870 }
2871
2872 let root = match object_bound::open_directory(quarantine) {
2876 Ok(root) => root,
2877 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2878 return Err(unsafe_path(QUARANTINE_VANISHED_REASON));
2879 }
2880 Err(source) => return Err(io_error(source)),
2881 };
2882 let root_identity = object_bound::identity_of(&root).map_err(io_error)?;
2883 if root_identity.device != expected.device || root_identity.inode != expected.inode {
2884 return Err(unsafe_path(
2885 "quarantined directory identity changed before deletion",
2886 ));
2887 }
2888
2889 let mut directories = vec![CapturedDirectory {
2891 file: root,
2892 name: CString::default(),
2893 parent: 0,
2894 depth: 0,
2895 identity: root_identity,
2896 }];
2897 let mut entries: Vec<CapturedEntry> = Vec::new();
2898 let mut pending = vec![0usize];
2899 while let Some(index) = pending.pop() {
2900 let depth = directories[index].depth;
2901 if depth > MAX_REMOVAL_DEPTH {
2902 return Err(unsafe_path(
2903 "managed directory nesting exceeds the removal depth cap",
2904 ));
2905 }
2906 let names = object_bound::list_names(&directories[index].file).map_err(io_error)?;
2907 for name in names {
2908 match object_bound::open_child_directory(&directories[index].file, &name) {
2909 Ok(child) => {
2910 let identity = object_bound::identity_of(&child).map_err(io_error)?;
2911 if !object_bound::same_filesystem(&root_identity, &identity) {
2912 return Err(unsafe_path(
2913 "managed directory crosses a filesystem boundary",
2914 ));
2915 }
2916 directories.push(CapturedDirectory {
2917 file: child,
2918 name,
2919 parent: index,
2920 depth: depth + 1,
2921 identity,
2922 });
2923 pending.push(directories.len() - 1);
2924 }
2925 Err(error) if object_bound::is_not_a_directory(&error) => {
2926 let (kind, identity) =
2927 object_bound::stat_entry(&directories[index].file, &name)
2928 .map_err(io_error)?;
2929 match kind {
2930 EntryKind::RegularFile | EntryKind::Symlink => {
2931 entries.push(CapturedEntry {
2932 parent: index,
2933 name,
2934 identity,
2935 });
2936 }
2937 EntryKind::Directory | EntryKind::Other => {
2938 return Err(unsafe_path(
2939 "managed directory contains an entry CAR cannot classify",
2940 ));
2941 }
2942 }
2943 }
2944 Err(error) if error.raw_os_error() == Some(libc::EXDEV) => {
2945 return Err(unsafe_path(
2946 "managed directory crosses a filesystem boundary",
2947 ));
2948 }
2949 Err(error)
2950 if matches!(
2951 error.raw_os_error(),
2952 Some(libc::ENOSYS) | Some(libc::EINVAL) | Some(libc::E2BIG)
2953 ) =>
2954 {
2955 return Err(unsafe_path(
2956 "descriptor-bound directory open is unavailable on this kernel",
2957 ));
2958 }
2959 Err(source) => return Err(io_error(source)),
2960 }
2961 }
2962 }
2963
2964 #[cfg(all(test, any(target_os = "macos", target_os = "linux")))]
2965 if let Some(hook) = &self.removal_hook {
2966 (hook.0)(RemovalPhase::AfterCapture, quarantine);
2967 }
2968
2969 for entry in &entries {
2972 let parent = &directories[entry.parent].file;
2973 match object_bound::stat_entry(parent, &entry.name) {
2974 Ok((_, identity)) if identity == entry.identity => {
2975 object_bound::unlink_entry(parent, &entry.name).map_err(io_error)?;
2976 }
2977 Ok(_) => {
2978 return Err(unsafe_path(
2979 "managed directory entry identity changed during deletion",
2980 ));
2981 }
2982 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
2983 Err(source) => return Err(io_error(source)),
2984 }
2985 }
2986
2987 while directories.len() > 1 {
2990 let CapturedDirectory {
2991 file,
2992 name,
2993 parent,
2994 identity,
2995 ..
2996 } = directories.pop().expect("non-empty");
2997 drop(file);
2998 let parent = &directories[parent].file;
2999 match object_bound::stat_entry(parent, &name) {
3000 Ok((EntryKind::Directory, actual)) if actual == identity => {
3001 object_bound::remove_empty_directory(parent, &name).map_err(|error| {
3002 if object_bound::is_not_empty(&error) {
3003 unsafe_path("managed directory contains entries that were not captured")
3004 } else {
3005 io_error(error)
3006 }
3007 })?;
3008 }
3009 Ok(_) => {
3010 return Err(unsafe_path(
3011 "managed directory entry identity changed during deletion",
3012 ));
3013 }
3014 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
3015 Err(source) => return Err(io_error(source)),
3016 }
3017 }
3018
3019 let CapturedDirectory { file: root, .. } = directories.pop().expect("root");
3025 drop(root);
3026 let container = quarantine
3027 .parent()
3028 .ok_or_else(|| unsafe_path("quarantine path has no parent"))?;
3029 let models =
3030 object_bound::open_directory(container).map_err(|source| ModelManagementError::Io {
3031 path: container.to_path_buf(),
3032 source,
3033 })?;
3034 let leaf = quarantine
3035 .file_name()
3036 .map(Path::new)
3037 .ok_or_else(|| unsafe_path("quarantine path has no file name"))?;
3038 let leaf = object_bound::c_name(leaf).map_err(io_error)?;
3039 match object_bound::stat_entry(&models, &leaf) {
3040 Ok((EntryKind::Directory, actual)) if actual == root_identity => {
3041 object_bound::remove_empty_directory(&models, &leaf).map_err(|error| {
3042 if object_bound::is_not_empty(&error) {
3043 unsafe_path("managed directory contains entries that were not captured")
3044 } else {
3045 io_error(error)
3046 }
3047 })
3048 }
3049 Ok(_) => Err(unsafe_path(
3050 "quarantined directory identity changed during deletion",
3051 )),
3052 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3055 Err(unsafe_path(QUARANTINE_VANISHED_REASON))
3056 }
3057 Err(source) => Err(io_error(source)),
3058 }
3059 }
3060}
3061
3062#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3064pub(crate) enum UnlinkKind {
3065 File,
3067 Link,
3069 FileOrLink,
3071}
3072
3073#[cfg(any(target_os = "macos", target_os = "linux"))]
3078pub(crate) fn unlink_checked(path: &Path, kind: UnlinkKind) -> std::io::Result<()> {
3079 use object_bound::EntryKind;
3080 let parent = path.parent().ok_or_else(|| {
3081 std::io::Error::new(std::io::ErrorKind::InvalidInput, "path has no parent")
3082 })?;
3083 let name = path.file_name().ok_or_else(|| {
3084 std::io::Error::new(std::io::ErrorKind::InvalidInput, "path has no file name")
3085 })?;
3086 let dir = object_bound::open_directory(parent)?;
3087 let name = object_bound::c_name(Path::new(name))?;
3088 let (actual, _) = object_bound::stat_entry(&dir, &name)?;
3089 let allowed = matches!(
3090 (kind, actual),
3091 (
3092 UnlinkKind::File | UnlinkKind::FileOrLink,
3093 EntryKind::RegularFile
3094 ) | (
3095 UnlinkKind::Link | UnlinkKind::FileOrLink,
3096 EntryKind::Symlink
3097 )
3098 );
3099 if !allowed {
3100 return Err(std::io::Error::new(
3101 std::io::ErrorKind::InvalidInput,
3102 format!("{} is no longer a {kind:?} entry", path.display()),
3103 ));
3104 }
3105 object_bound::unlink_entry(&dir, &name)
3106}
3107
3108#[cfg(not(any(target_os = "macos", target_os = "linux")))]
3111pub(crate) fn unlink_checked(path: &Path, kind: UnlinkKind) -> std::io::Result<()> {
3112 let meta = std::fs::symlink_metadata(path)?;
3113 let allowed = match kind {
3114 UnlinkKind::File => meta.is_file(),
3115 UnlinkKind::Link => meta.file_type().is_symlink(),
3116 UnlinkKind::FileOrLink => meta.is_file() || meta.file_type().is_symlink(),
3117 };
3118 if !allowed {
3119 return Err(std::io::Error::new(
3120 std::io::ErrorKind::InvalidInput,
3121 format!("{} is no longer a {kind:?} entry", path.display()),
3122 ));
3123 }
3124 std::fs::remove_file(path)
3125}
3126
3127#[cfg(not(unix))]
3132fn sync_directory(_path: &Path) -> std::io::Result<()> {
3133 Ok(())
3134}
3135
3136#[cfg(test)]
3137mod tests {
3138 #[cfg(any(target_os = "macos", target_os = "linux"))]
3141 #[test]
3142 fn interrupted_hub_retirements_resume() {
3143 let root = tempfile::tempdir().unwrap();
3144 let store =
3145 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3146 let hub = root.path().join("hub");
3147 let repo = hub.join("models--a--b");
3148 std::fs::create_dir_all(repo.join("blobs")).unwrap();
3149 std::fs::write(repo.join("blobs/x"), b"bytes").unwrap();
3150 let quarantine = hub.join(".car-retire-1");
3151 let journal = DetachedTreeJournal {
3152 source: repo.clone(),
3153 quarantine: quarantine.clone(),
3154 identity: artifact_identity(&repo, "t").unwrap(),
3155 };
3156 let journals = store.detached_tree_journals();
3157 std::fs::create_dir_all(&journals).unwrap();
3158 write_private_json(&journals.join("j1"), &journal).unwrap();
3159 std::fs::rename(&repo, &quarantine).unwrap();
3160 let never = DetachedTreeJournal {
3162 source: hub.join("models--c--d"),
3163 quarantine: hub.join(".car-retire-2"),
3164 identity: journal.identity.clone(),
3165 };
3166 write_private_json(&journals.join("j2"), &never).unwrap();
3167
3168 assert!(store.resume_detached_trees().is_empty());
3169 assert!(!quarantine.exists());
3170 assert!(std::fs::read_dir(&journals).unwrap().next().is_none());
3171
3172 let foreign = hub.join("models--e--f");
3175 std::fs::create_dir_all(&foreign).unwrap();
3176 std::fs::write(foreign.join("keep"), b"k").unwrap();
3177 let aimed = DetachedTreeJournal {
3178 source: hub.join("models--g--h"),
3179 quarantine: foreign.clone(),
3180 identity: artifact_identity(&foreign, "t").unwrap(),
3181 };
3182 write_private_json(&journals.join("j3"), &aimed).unwrap();
3183 let swapped = hub.join(".car-retire-3");
3185 std::fs::create_dir_all(&swapped).unwrap();
3186 std::fs::write(swapped.join("other"), b"o").unwrap();
3187 let stale = DetachedTreeJournal {
3188 source: hub.join("models--i--j"),
3189 quarantine: swapped.clone(),
3190 identity: journal.identity.clone(),
3191 };
3192 write_private_json(&journals.join("j4"), &stale).unwrap();
3193
3194 let errors = store.resume_detached_trees();
3195 assert_eq!(errors.len(), 1, "{errors:?}");
3196 assert!(foreign.join("keep").exists());
3197 assert!(swapped.join("other").exists());
3198 assert!(journals.join("j3").exists() && journals.join("j4").exists());
3199 }
3200
3201 #[test]
3205 fn usage_tracking_starts_with_the_store_and_stays_put() {
3206 let root = tempfile::tempdir().unwrap();
3207 let store =
3208 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3209 assert_eq!(
3210 store.usage_tracking_since(),
3211 None,
3212 "construction starts nothing"
3213 );
3214 store.start_usage_tracking();
3215 let since = store
3216 .usage_tracking_since()
3217 .expect("marker written when started");
3218 let marker = store.usage_root.join(USAGE_TRACKING_MARKER);
3219 write_private_json(&marker, &UsageTracking { since: 7 }).unwrap();
3220 let again =
3221 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3222 again.start_usage_tracking();
3223 assert_eq!(again.usage_tracking_since(), Some(7));
3224 assert!(since > 7);
3225 assert_eq!(store.last_used("org/model"), None, "never used is not used");
3226 }
3227
3228 #[test]
3230 fn a_lease_stamps_use_at_most_once_per_interval() {
3231 let root = tempfile::tempdir().unwrap();
3232 let store =
3233 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3234 drop(store.acquire_lease("org/model").unwrap());
3235 let stamped = store.last_used("org/model").expect("first lease stamps");
3236 assert!(stamped > 0);
3237 std::fs::remove_file(store.usage_root.join(hashed_filename("org/model"))).unwrap();
3238 drop(store.acquire_lease("org/model").unwrap());
3239 assert_eq!(
3240 store.last_used("org/model"),
3241 None,
3242 "throttled: no second write"
3243 );
3244 drop(store.acquire_lease("org/other").unwrap());
3246 assert!(store.last_used("org/other").is_some());
3247 write_private_json(
3249 &store.usage_root.join(hashed_filename("org/model")),
3250 &UsageStamp {
3251 model_id: "org/forged".into(),
3252 last_used: 1,
3253 },
3254 )
3255 .unwrap();
3256 assert_eq!(store.last_used("org/model"), None);
3257 }
3258
3259 #[test]
3262 fn coverage_runs_continue_across_restarts_and_break_on_long_gaps() {
3263 let root = tempfile::tempdir().unwrap();
3264 let store =
3265 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3266 assert_eq!(store.usage_covered_since(), None);
3267 store.beat();
3268 let first = store.usage_covered_since().unwrap();
3269 store.beat();
3270 assert_eq!(store.usage_covered_since(), Some(first), "same run");
3271 let restarted =
3273 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3274 restarted.beat();
3275 assert_eq!(restarted.usage_covered_since(), Some(first));
3276 write_private_json(
3278 &restarted.usage_root.join(USAGE_COVERAGE),
3279 &UsageCoverage {
3280 run_start: 5,
3281 last_beat: 10,
3282 },
3283 )
3284 .unwrap();
3285 let later =
3286 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3287 later.beat();
3288 assert!(later.usage_covered_since().unwrap() > 10);
3289 }
3290
3291 #[test]
3292 fn an_interrupted_retirement_clears_its_tombstone_at_resume() {
3293 let root = tempfile::tempdir().unwrap();
3294 let store =
3295 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3296 let id = "mlx/some:4bit";
3297 store.note_retiring(id).unwrap();
3298 write_private_json(
3299 &store.tombstone_path(id),
3300 &RemovalTombstone {
3301 model_id: id.into(),
3302 removal_generation: 1,
3303 artifact_kind: None,
3304 preserved_shared_cache_references: Vec::new(),
3305 },
3306 )
3307 .unwrap();
3308 assert!(!store.car_enabled(id).unwrap());
3309 assert!(store.tombstone_path(id).exists());
3310 assert!(store.resume_retiring().is_empty());
3311 assert!(!store.tombstone_path(id).exists());
3312 assert!(!store.retiring_intent_path(id).exists());
3313 }
3314
3315 use super::*;
3316
3317 #[test]
3319 fn identity_initialization_child() {
3320 use std::io::{BufRead, Write};
3321 let Some(root) = std::env::var_os("CAR_IDENTITY_TEST_ROOT") else {
3322 return;
3323 };
3324 let root = PathBuf::from(root);
3325 let mut store = ModelManagementStore::new(root.join("state"), root.join("models"));
3326 if std::env::var("CAR_IDENTITY_TEST_ROLE").as_deref() == Ok("creator") {
3327 store.identity_init_hook = Some(|path| {
3328 let phase = std::env::var("CAR_IDENTITY_TEST_PHASE").unwrap();
3329 if path.parent().unwrap().file_name().unwrap() == phase.as_str() {
3330 println!("IDENTITY_CREATED");
3331 std::io::stdout().flush().unwrap();
3332 let mut release = String::new();
3333 std::io::stdin().lock().read_line(&mut release).unwrap();
3334 assert_eq!(release.trim(), "release");
3335 }
3336 });
3337 } else {
3338 println!("READER_STARTED");
3339 std::io::stdout().flush().unwrap();
3340 }
3341 let _lease = store.acquire_lease("fixture/model").unwrap();
3342 }
3343
3344 #[test]
3345 fn identity_initialization_serializes_processes_before_json_is_written() {
3346 use std::io::{BufRead, BufReader, Write};
3347 use std::process::{Command, Stdio};
3348 fn wait_marker(reader: &mut impl BufRead, marker: &str) {
3349 let mut line = String::new();
3350 loop {
3351 line.clear();
3352 assert_ne!(
3353 reader.read_line(&mut line).unwrap(),
3354 0,
3355 "child exited before {marker}"
3356 );
3357 if line.contains(marker) {
3358 return;
3359 }
3360 }
3361 }
3362 for phase in ["mutation", "activity"] {
3363 let root = tempfile::tempdir().unwrap();
3364 let spawn = |role: &str| {
3365 Command::new(std::env::current_exe().unwrap())
3366 .args([
3367 "--exact",
3368 "model_management::tests::identity_initialization_child",
3369 "--nocapture",
3370 ])
3371 .env("CAR_IDENTITY_TEST_ROOT", root.path())
3372 .env("CAR_IDENTITY_TEST_ROLE", role)
3373 .env("CAR_IDENTITY_TEST_PHASE", phase)
3374 .stdin(Stdio::piped())
3375 .stdout(Stdio::piped())
3376 .stderr(Stdio::inherit())
3377 .spawn()
3378 .unwrap()
3379 };
3380 let mut creator = spawn("creator");
3381 let mut creator_out = BufReader::new(creator.stdout.take().unwrap());
3382 wait_marker(&mut creator_out, "IDENTITY_CREATED");
3383 let store =
3384 ModelManagementStore::new(root.path().join("state"), root.path().join("models"));
3385 let path = if phase == "mutation" {
3386 store.mutation_lock_path("fixture/model")
3387 } else {
3388 store.lease_path("fixture/model")
3389 };
3390 assert!(
3391 std::fs::read(&path).unwrap().is_empty(),
3392 "probe must hit the original partial-publication window"
3393 );
3394 let gate =
3395 open_existing_identity_lock(&path.parent().unwrap().join(".initialization.lock"))
3396 .unwrap();
3397 assert!(
3398 gate.try_lock().is_err(),
3399 "creator must hold the initialization gate before publishing an empty identity"
3400 );
3401 let mut contender = spawn("reader");
3402 let mut contender_out = BufReader::new(contender.stdout.take().unwrap());
3403 wait_marker(&mut contender_out, "READER_STARTED");
3404 assert!(contender.try_wait().unwrap().is_none());
3405 writeln!(creator.stdin.take().unwrap(), "release").unwrap();
3406 assert!(creator.wait().unwrap().success());
3407 assert!(contender.wait().unwrap().success());
3408 let before = std::fs::metadata(&path).unwrap();
3409 drop(store.acquire_lease("fixture/model").unwrap());
3410 #[cfg(unix)]
3411 {
3412 use std::os::unix::fs::MetadataExt;
3413 assert_eq!(before.ino(), std::fs::metadata(&path).unwrap().ino());
3414 }
3415 #[cfg(not(unix))]
3416 let _ = before;
3417 }
3418 }
3419
3420 #[test]
3421 fn identity_initialization_preserves_persistent_corruption_refusal() {
3422 for phase in ["mutation", "activity"] {
3423 for corrupt in [b"".as_slice(), b"not json".as_slice()] {
3424 let root = tempfile::tempdir().unwrap();
3425 let store = ModelManagementStore::new(
3426 root.path().join("state"),
3427 root.path().join("models"),
3428 );
3429 drop(store.acquire_lease("fixture/model").unwrap());
3430 let path = if phase == "mutation" {
3431 store.mutation_lock_path("fixture/model")
3432 } else {
3433 store.lease_path("fixture/model")
3434 };
3435 std::fs::write(&path, corrupt).unwrap();
3436 assert!(matches!(
3437 store.acquire_lease("fixture/model"),
3438 Err(ModelManagementError::InvalidState { .. })
3439 ));
3440 assert_eq!(std::fs::read(&path).unwrap(), corrupt);
3441 }
3442 }
3443 }
3444
3445 #[cfg(unix)]
3446 #[test]
3447 fn identity_initialization_rejects_linked_gates_and_identity_files() {
3448 for phase in ["mutation", "activity"] {
3449 for gate in [false, true] {
3450 for symlink in [false, true] {
3451 let root = tempfile::tempdir().unwrap();
3452 let store = ModelManagementStore::new(
3453 root.path().join("state"),
3454 root.path().join("models"),
3455 );
3456 let identity = if phase == "mutation" {
3457 store.mutation_lock_path("fixture/model")
3458 } else {
3459 store.lease_path("fixture/model")
3460 };
3461 create_private_dir(identity.parent().unwrap()).unwrap();
3462 let path = if gate {
3463 identity.parent().unwrap().join(".initialization.lock")
3464 } else {
3465 identity
3466 };
3467 let outside = root.path().join("outside");
3468 std::fs::write(&outside, b"untouched").unwrap();
3469 if symlink {
3470 std::os::unix::fs::symlink(&outside, &path).unwrap();
3471 } else {
3472 std::fs::hard_link(&outside, &path).unwrap();
3473 }
3474 assert!(matches!(
3475 store.acquire_lease("fixture/model"),
3476 Err(ModelManagementError::InvalidState { .. })
3477 ));
3478 assert_eq!(std::fs::read(&outside).unwrap(), b"untouched");
3479 }
3480 }
3481 }
3482 }
3483
3484 fn directory_receipt(model_id: &str, managed_path: PathBuf) -> InstallReceipt {
3485 InstallReceipt {
3486 model_id: model_id.into(),
3487 managed_path,
3488 artifact_kind: ManagedArtifactKind::Directory,
3489 source_model_id: model_id.into(),
3490 source_revision: None,
3491 creation_generation: 1,
3492 shared_cache_references: vec![],
3493 adopted: false,
3494 }
3495 }
3496
3497 #[test]
3498 fn raw_receipt_identity_is_private_and_collision_fails_closed() {
3499 let root = tempfile::tempdir().unwrap();
3500 let models = root.path().join("models");
3501 std::fs::create_dir_all(&models).unwrap();
3502 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3503 let receipt = directory_receipt("org/model", models.join("Model"));
3504 store.write_receipt(&receipt).unwrap();
3505 assert_eq!(store.load_receipt("org/model").unwrap(), Some(receipt));
3506
3507 #[cfg(unix)]
3508 {
3509 use std::os::unix::fs::PermissionsExt;
3510 assert_eq!(
3511 std::fs::metadata(store.receipt_path("org/model"))
3512 .unwrap()
3513 .permissions()
3514 .mode()
3515 & 0o077,
3516 0
3517 );
3518 }
3519
3520 let collision = directory_receipt("different/model", models.join("Other"));
3521 std::fs::write(
3522 store.receipt_path("org/model"),
3523 serde_json::to_vec(&collision).unwrap(),
3524 )
3525 .unwrap();
3526 assert!(matches!(
3527 store.load_receipt("org/model"),
3528 Err(ModelManagementError::IdentityMismatch { .. })
3529 ));
3530 }
3531
3532 #[test]
3533 fn journaled_removal_rejects_leaf_identity_swap() {
3534 let root = tempfile::tempdir().unwrap();
3535 let models = root.path().join("models");
3536 std::fs::create_dir_all(&models).unwrap();
3537 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3538 let managed = models.join("Managed");
3539 std::fs::create_dir_all(&managed).unwrap();
3540 std::fs::write(managed.join("weights"), b"original").unwrap();
3541 let receipt = directory_receipt("org/model", managed.clone());
3542 store.write_receipt(&receipt).unwrap();
3543 let journal = RemovalJournal {
3544 model_id: "org/model".into(),
3545 removal_generation: 2,
3546 artifact_identity: artifact_identity(&managed, "org/model").unwrap(),
3547 receipt,
3548 quarantine_path: store.allocate_quarantine("org/model").unwrap(),
3549 };
3550 write_private_json(&store.removal_journal_path("org/model"), &journal).unwrap();
3551 let displaced = models.join("displaced-original");
3552 std::fs::rename(&managed, &displaced).unwrap();
3553 std::fs::create_dir_all(&managed).unwrap();
3554 std::fs::write(managed.join("sentinel"), b"replacement").unwrap();
3555
3556 assert!(matches!(
3557 store.begin_mutation("org/model").unwrap().remove(2),
3558 Err(ModelManagementError::UnsafeManagedPath { .. })
3559 ));
3560 assert_eq!(
3561 std::fs::read(managed.join("sentinel")).unwrap(),
3562 b"replacement"
3563 );
3564 }
3565
3566 #[test]
3567 fn journaled_removal_rejects_models_parent_swap() {
3568 let root = tempfile::tempdir().unwrap();
3569 let models = root.path().join("models");
3570 std::fs::create_dir_all(&models).unwrap();
3571 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3572 let managed = models.join("Managed");
3573 std::fs::create_dir_all(&managed).unwrap();
3574 std::fs::write(managed.join("weights"), b"original").unwrap();
3575 let receipt = directory_receipt("org/model", managed.clone());
3576 store.write_receipt(&receipt).unwrap();
3577 let journal = RemovalJournal {
3578 model_id: "org/model".into(),
3579 removal_generation: 2,
3580 artifact_identity: artifact_identity(&managed, "org/model").unwrap(),
3581 receipt,
3582 quarantine_path: store.allocate_quarantine("org/model").unwrap(),
3583 };
3584 write_private_json(&store.removal_journal_path("org/model"), &journal).unwrap();
3585 std::fs::rename(&models, root.path().join("models-original")).unwrap();
3586 std::fs::create_dir_all(&managed).unwrap();
3587 std::fs::write(managed.join("sentinel"), b"replacement-parent").unwrap();
3588
3589 assert!(matches!(
3590 store.begin_mutation("org/model").unwrap().remove(2),
3591 Err(ModelManagementError::UnsafeManagedPath { .. })
3592 ));
3593 assert_eq!(
3594 std::fs::read(managed.join("sentinel")).unwrap(),
3595 b"replacement-parent"
3596 );
3597 }
3598
3599 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
3600 #[test]
3601 fn directory_receipts_fail_closed_without_recursive_cleanup() {
3602 let root = tempfile::tempdir().unwrap();
3603 let models = root.path().join("models");
3604 std::fs::create_dir_all(&models).unwrap();
3605 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3606 let managed = models.join("Managed");
3607 std::fs::create_dir_all(&managed).unwrap();
3608 std::fs::write(managed.join("shared-sentinel"), b"preserve").unwrap();
3609 store
3610 .write_receipt(&directory_receipt("org/model", managed.clone()))
3611 .unwrap();
3612
3613 assert!(!store.can_remove("org/model").unwrap());
3614 let error = store
3615 .begin_mutation("org/model")
3616 .unwrap()
3617 .remove(2)
3618 .unwrap_err();
3619 assert!(matches!(
3620 error,
3621 ModelManagementError::UnsafeManagedPath { .. }
3622 ));
3623 assert_eq!(
3624 std::fs::read(managed.join("shared-sentinel")).unwrap(),
3625 b"preserve"
3626 );
3627 }
3628
3629 #[test]
3630 fn failed_directory_staging_is_left_for_object_bound_cleanup() {
3631 let root = tempfile::tempdir().unwrap();
3632 let models = root.path().join("models");
3633 std::fs::create_dir_all(&models).unwrap();
3634 let store = ModelManagementStore::new(root.path().join("state"), models);
3635 let staging = store.create_install_staging("org/model").unwrap();
3636 std::fs::write(staging.join("shared-sentinel"), b"preserve").unwrap();
3637
3638 store.discard_install_staging(&staging);
3639
3640 assert_eq!(
3641 std::fs::read(staging.join("shared-sentinel")).unwrap(),
3642 b"preserve"
3643 );
3644 }
3645
3646 #[cfg(unix)]
3647 #[test]
3648 fn completed_symlink_removal_retry_is_idempotent() {
3649 use std::os::unix::fs::symlink;
3650
3651 let root = tempfile::tempdir().unwrap();
3652 let models = root.path().join("models");
3653 let shared = root.path().join("shared");
3654 std::fs::create_dir_all(&models).unwrap();
3655 std::fs::create_dir_all(&shared).unwrap();
3656 std::fs::write(shared.join("sentinel"), b"preserve").unwrap();
3657 let managed = models.join("Managed");
3658 symlink(&shared, &managed).unwrap();
3659 let store = ModelManagementStore::new(root.path().join("state"), models);
3660 let receipt = InstallReceipt {
3661 model_id: "org/model".into(),
3662 managed_path: managed,
3663 artifact_kind: ManagedArtifactKind::Symlink,
3664 source_model_id: "org/model".into(),
3665 source_revision: None,
3666 creation_generation: 1,
3667 shared_cache_references: vec![shared.canonicalize().unwrap()],
3668 adopted: false,
3669 };
3670 store.write_receipt(&receipt).unwrap();
3671
3672 let first = store
3673 .begin_mutation("org/model")
3674 .unwrap()
3675 .remove(2)
3676 .unwrap();
3677 let retry = store
3678 .begin_mutation("org/model")
3679 .unwrap()
3680 .remove(3)
3681 .unwrap();
3682 assert_eq!(retry, first);
3683 assert!(shared.join("sentinel").exists());
3684 }
3685
3686 #[test]
3687 fn pulled_directory_publish_before_receipt_resumes_from_install_intent() {
3688 let root = tempfile::tempdir().unwrap();
3689 let models = root.path().join("models");
3690 std::fs::create_dir_all(&models).unwrap();
3691 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3692 let staging = store.create_install_staging("org/model").unwrap();
3693 std::fs::write(staging.join("weights"), b"owned").unwrap();
3694 let receipt = store
3695 .install_receipt_for_publication(
3696 "org/model",
3697 "org/model",
3698 1,
3699 false,
3700 models.join("Managed"),
3701 ManagedArtifactKind::Directory,
3702 &staging,
3703 )
3704 .unwrap();
3705 store.begin_install_intent(receipt, Some(&staging)).unwrap();
3706 store
3707 .publish_install_staging("org/model", &staging, "Managed")
3708 .unwrap();
3709 assert!(store.load_receipt("org/model").unwrap().is_none());
3710
3711 let resumed = store.resume_install_intent("org/model").unwrap().unwrap();
3712 assert_eq!(resumed.managed_path, models.join("Managed"));
3713 assert!(store.load_receipt("org/model").unwrap().is_some());
3714 #[cfg(any(target_os = "macos", target_os = "linux"))]
3715 assert!(store.can_remove("org/model").unwrap());
3716 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
3717 assert!(!store.can_remove("org/model").unwrap());
3718 }
3719
3720 #[cfg(unix)]
3721 #[test]
3722 fn adopted_projection_publish_before_receipt_resumes_from_install_intent() {
3723 let root = tempfile::tempdir().unwrap();
3724 let models = root.path().join("models");
3725 let source = root.path().join("hand-installed");
3726 std::fs::create_dir_all(&models).unwrap();
3727 std::fs::create_dir_all(&source).unwrap();
3728 std::fs::write(source.join("sentinel"), b"preserve").unwrap();
3729 let store = ModelManagementStore::new(root.path().join("state"), models);
3730 let managed = store.adopted_projection_path("org/model").unwrap();
3731 let receipt = store
3732 .install_receipt_for_publication(
3733 "org/model",
3734 "org/model",
3735 1,
3736 true,
3737 managed.clone(),
3738 ManagedArtifactKind::Symlink,
3739 &source,
3740 )
3741 .unwrap();
3742 store.begin_install_intent(receipt, None).unwrap();
3743 store
3744 .materialize_adopted_projection("org/model", &source)
3745 .unwrap();
3746 assert!(store.load_receipt("org/model").unwrap().is_none());
3747
3748 let resumed = store.resume_install_intent("org/model").unwrap().unwrap();
3749 assert_eq!(resumed.managed_path, managed);
3750 assert!(store.can_remove("org/model").unwrap());
3751 assert!(source.join("sentinel").exists());
3752 }
3753
3754 #[cfg(unix)]
3755 #[test]
3756 fn install_intent_never_claims_a_preexisting_unreceipted_leaf() {
3757 use std::os::unix::fs::symlink;
3758
3759 let root = tempfile::tempdir().unwrap();
3760 let models = root.path().join("models");
3761 let source = root.path().join("shared");
3762 std::fs::create_dir_all(&models).unwrap();
3763 std::fs::create_dir_all(&source).unwrap();
3764 let managed = models.join("Managed");
3765 symlink(&source, &managed).unwrap();
3766 let store = ModelManagementStore::new(root.path().join("state"), models);
3767 let receipt = store
3768 .install_receipt_for_publication(
3769 "org/model",
3770 "org/model",
3771 1,
3772 false,
3773 managed.clone(),
3774 ManagedArtifactKind::Symlink,
3775 &source,
3776 )
3777 .unwrap();
3778
3779 let error = store.begin_install_intent(receipt, None).unwrap_err();
3780 assert!(matches!(
3781 error,
3782 ModelManagementError::UnsafeManagedPath { .. }
3783 ));
3784 assert!(managed.exists());
3785 assert!(store.load_receipt("org/model").unwrap().is_none());
3786 }
3787
3788 #[test]
3789 fn forged_install_journal_paths_fail_before_filesystem_mutation() {
3790 let root = tempfile::tempdir().unwrap();
3791 let models = root.path().join("models");
3792 std::fs::create_dir_all(&models).unwrap();
3793 let outside = root.path().join("outside");
3794 std::fs::create_dir_all(&outside).unwrap();
3795 std::fs::write(outside.join("sentinel"), b"preserve").unwrap();
3796 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3797
3798 let forged_managed = InstallJournal {
3799 model_id: "org/model".into(),
3800 receipt: directory_receipt("org/model", outside.clone()),
3801 staging_path: None,
3802 staging_identity: None,
3803 };
3804 write_private_json(&store.install_journal_path("org/model"), &forged_managed).unwrap();
3805 assert!(matches!(
3806 store.resume_install_intent("org/model"),
3807 Err(ModelManagementError::UnsafeManagedPath { .. })
3808 ));
3809 assert!(outside.join("sentinel").exists());
3810
3811 let forged_staging = InstallJournal {
3812 model_id: "org/other".into(),
3813 receipt: directory_receipt("org/other", models.join("Other")),
3814 staging_path: Some(outside.clone()),
3815 staging_identity: Some(artifact_identity(&outside, "org/other").unwrap()),
3816 };
3817 write_private_json(&store.install_journal_path("org/other"), &forged_staging).unwrap();
3818 assert!(matches!(
3819 store.resume_install_intent("org/other"),
3820 Err(ModelManagementError::UnsafeManagedPath { .. })
3821 ));
3822 assert!(outside.join("sentinel").exists());
3823 }
3824
3825 #[cfg(unix)]
3826 #[test]
3827 fn startup_ignores_noncanonical_removal_journal_without_touching_valid_model() {
3828 use std::os::unix::fs::symlink;
3829
3830 let root = tempfile::tempdir().unwrap();
3831 let state = root.path().join("state");
3832 let models = root.path().join("models");
3833 let shared = root.path().join("shared");
3834 let outside = root.path().join("outside");
3835 std::fs::create_dir_all(&models).unwrap();
3836 std::fs::create_dir_all(&shared).unwrap();
3837 std::fs::create_dir_all(&outside).unwrap();
3838 std::fs::write(shared.join("sentinel"), b"shared").unwrap();
3839 std::fs::write(outside.join("sentinel"), b"outside").unwrap();
3840 let managed = models.join("Managed");
3841 symlink(&shared, &managed).unwrap();
3842 let store = ModelManagementStore::new(state.clone(), models.clone());
3843 let legitimate = InstallReceipt {
3844 model_id: "org/model".into(),
3845 managed_path: managed.clone(),
3846 artifact_kind: ManagedArtifactKind::Symlink,
3847 source_model_id: "org/model".into(),
3848 source_revision: None,
3849 creation_generation: 1,
3850 shared_cache_references: vec![shared.canonicalize().unwrap()],
3851 adopted: false,
3852 };
3853 store.write_receipt(&legitimate).unwrap();
3854 create_private_dir(&store.removals_root).unwrap();
3855 let forged = RemovalJournal {
3856 model_id: "org/model".into(),
3857 removal_generation: 99,
3858 receipt: InstallReceipt {
3859 managed_path: outside.join("missing-leaf"),
3860 ..legitimate
3861 },
3862 quarantine_path: models.join(".car-remove-forged"),
3863 artifact_identity: ArtifactIdentity {
3864 kind: ManagedArtifactKind::Symlink,
3865 device: 1,
3866 inode: 1,
3867 file_len: 0,
3868 },
3869 };
3870 write_private_json(&store.removals_root.join("junk.json"), &forged).unwrap();
3871 drop(store);
3872
3873 let recovered = ModelManagementStore::new(state, models);
3874 assert!(std::fs::symlink_metadata(&managed).is_ok());
3875 assert!(recovered.car_enabled("org/model").unwrap());
3876 assert_eq!(std::fs::read(shared.join("sentinel")).unwrap(), b"shared");
3877 assert_eq!(std::fs::read(outside.join("sentinel")).unwrap(), b"outside");
3878 }
3879
3880 #[cfg(any(target_os = "macos", target_os = "linux"))]
3883 #[test]
3884 fn mutation_guard_releases_lock_even_with_an_inherited_descriptor() {
3885 let (_root, _models, store) = supported_platform_store();
3886 let guard = store.begin_mutation("org/model").unwrap();
3887 let inherited = guard._file.try_clone().unwrap();
3888 assert!(matches!(
3889 store.begin_mutation("org/model"),
3890 Err(ModelManagementError::MutationInProgress { .. })
3891 ));
3892 drop(guard);
3893 let replacement = store.begin_mutation("org/model").unwrap();
3894 drop(inherited);
3896 assert!(matches!(
3897 store.begin_mutation("org/model"),
3898 Err(ModelManagementError::MutationInProgress { .. })
3899 ));
3900 drop(replacement);
3901 assert!(store.begin_mutation("org/model").is_ok());
3902 }
3903
3904 #[cfg(any(target_os = "macos", target_os = "linux"))]
3905 #[test]
3906 fn lease_and_idle_guards_release_despite_inherited_descriptors() {
3907 let (_root, _models, store) = supported_platform_store();
3908 let lease = store.acquire_lease("org/model").unwrap();
3909 let inherited = lease._file.try_clone().unwrap();
3910 assert!(store.model_in_use("org/model").unwrap());
3911 drop(lease);
3912 assert!(!store.model_in_use("org/model").unwrap());
3913 let idle = store.lock_idle_activity("org/model").unwrap();
3914 let inherited_idle = idle.try_clone().unwrap();
3915 assert!(store.model_in_use("org/model").unwrap());
3916 drop(idle);
3917 assert!(!store.model_in_use("org/model").unwrap());
3918 drop((inherited, inherited_idle));
3919 }
3920
3921 #[cfg(any(target_os = "macos", target_os = "linux"))]
3922 #[test]
3923 fn temporary_owner_releases_initialization_gate_despite_inherited_descriptor() {
3924 let (_root, _models, store) = supported_platform_store();
3925 let identity = store
3926 .open_identity_lock(&store.mutation_lock_path("org/model"), "org/model")
3927 .unwrap();
3928 let path = store
3929 .mutation_lock_path("org/model")
3930 .parent()
3931 .unwrap()
3932 .join(".initialization.lock");
3933 let gate = open_existing_identity_lock(&path).unwrap();
3934 gate.lock().unwrap();
3935 let gate = OwnedFileLock::new(gate);
3936 let inherited = gate.try_clone().unwrap();
3937 let contender = open_existing_identity_lock(&path).unwrap();
3938 assert!(contender.try_lock().is_err());
3939 drop(gate);
3940 contender.try_lock().unwrap();
3941 contender.unlock().unwrap();
3942 drop((identity, inherited));
3943 }
3944
3945 #[cfg(any(target_os = "macos", target_os = "linux"))]
3946 fn supported_platform_store() -> (tempfile::TempDir, PathBuf, ModelManagementStore) {
3947 let root = tempfile::tempdir().unwrap();
3948 let models = root.path().join("models");
3949 std::fs::create_dir_all(&models).unwrap();
3950 let store = ModelManagementStore::new(root.path().join("state"), models.clone());
3951 (root, models, store)
3952 }
3953
3954 #[cfg(any(target_os = "macos", target_os = "linux"))]
3955 fn no_quarantine_left(models: &Path) -> bool {
3956 std::fs::read_dir(models).unwrap().all(|entry| {
3957 !entry
3958 .unwrap()
3959 .file_name()
3960 .to_string_lossy()
3961 .starts_with(".car-remove-")
3962 })
3963 }
3964
3965 #[cfg(any(target_os = "macos", target_os = "linux"))]
3966 const MAX_REMOVAL_DEPTH_FOR_TESTS: usize = super::object_bound::MAX_REMOVAL_DEPTH;
3967
3968 #[cfg(any(target_os = "macos", target_os = "linux"))]
3972 fn unsafe_managed_reason(error: &ModelManagementError) -> &str {
3973 match error {
3974 ModelManagementError::UnsafeManagedPath { reason, .. } => reason,
3975 other => panic!("expected UnsafeManagedPath, got {other:?}"),
3976 }
3977 }
3978
3979 #[cfg(any(target_os = "macos", target_os = "linux"))]
3980 #[test]
3981 fn directory_receipt_is_removable_and_leaves_a_tombstone() {
3982 let (_root, models, store) = supported_platform_store();
3983 let managed = models.join("Managed");
3984 std::fs::create_dir_all(&managed).unwrap();
3985 std::fs::write(managed.join("weights"), b"owned").unwrap();
3986 store
3987 .write_receipt(&directory_receipt("org/model", managed.clone()))
3988 .unwrap();
3989
3990 assert!(store.can_remove("org/model").unwrap());
3991 let result = store
3992 .begin_mutation("org/model")
3993 .unwrap()
3994 .remove(2)
3995 .unwrap();
3996
3997 assert_eq!(result.artifact_kind, ManagedArtifactKind::Directory);
3998 assert!(std::fs::symlink_metadata(&managed).is_err());
3999 assert!(store.load_receipt("org/model").unwrap().is_none());
4000 let tombstone = store.load_tombstone("org/model").unwrap().unwrap();
4001 assert_eq!(
4002 tombstone.artifact_kind,
4003 Some(ManagedArtifactKind::Directory)
4004 );
4005 assert_eq!(tombstone.removal_generation, 2);
4006 assert!(no_quarantine_left(&models));
4007 assert!(store.load_removal_journal("org/model").unwrap().is_none());
4008 }
4009
4010 #[cfg(any(target_os = "macos", target_os = "linux"))]
4011 #[test]
4012 fn directory_removal_preserves_shared_cache_symlink_targets() {
4013 use std::os::unix::fs::symlink;
4014
4015 let (root, models, store) = supported_platform_store();
4016 let shared = root.path().join("shared-cache");
4017 std::fs::create_dir_all(&shared).unwrap();
4018 std::fs::write(shared.join("blob"), b"shared-bytes").unwrap();
4019 let managed = models.join("Managed");
4020 let nested = managed.join("a").join("b").join("c");
4021 std::fs::create_dir_all(&nested).unwrap();
4022 std::fs::write(managed.join("config.json"), b"{}").unwrap();
4023 std::fs::write(nested.join("weights"), b"owned").unwrap();
4024 symlink(shared.join("blob"), managed.join("blob.safetensors")).unwrap();
4025 let blob = shared.join("blob").canonicalize().unwrap();
4026 let mut receipt = directory_receipt("org/model", managed.clone());
4027 receipt.shared_cache_references = vec![blob.clone()];
4028 store.write_receipt(&receipt).unwrap();
4029
4030 let result = store
4031 .begin_mutation("org/model")
4032 .unwrap()
4033 .remove(2)
4034 .unwrap();
4035
4036 assert!(std::fs::symlink_metadata(&managed).is_err());
4037 assert_eq!(std::fs::read(&blob).unwrap(), b"shared-bytes");
4038 assert_eq!(std::fs::read_dir(&shared).unwrap().count(), 1);
4039 assert_eq!(result.preserved_shared_cache_references, vec![blob]);
4040 assert!(no_quarantine_left(&models));
4041 }
4042
4043 #[cfg(any(target_os = "macos", target_os = "linux"))]
4048 #[test]
4049 fn the_revoked_copies_report_never_resumes_a_removal() {
4050 let (root, models, store) = supported_platform_store();
4051 let state = root.path().join("state");
4052 let managed = models.join("Managed");
4053 std::fs::create_dir_all(&managed).unwrap();
4054 std::fs::write(managed.join("weights"), b"owned").unwrap();
4055 let receipt = directory_receipt("org/model", managed.clone());
4056 store.write_receipt(&receipt).unwrap();
4057 let quarantine = store.allocate_quarantine("org/model").unwrap();
4058 let journal = RemovalJournal {
4059 model_id: "org/model".into(),
4060 removal_generation: 2,
4061 artifact_identity: artifact_identity(&managed, "org/model").unwrap(),
4062 receipt: receipt.clone(),
4063 quarantine_path: quarantine.clone(),
4064 };
4065 write_private_json(&store.removal_journal_path("org/model"), &journal).unwrap();
4066 std::fs::rename(&managed, &quarantine).unwrap();
4067 drop(store);
4068 crate::catalog::save_revoked_names(
4069 &crate::catalog::revoked_names_path(&state),
4070 &std::collections::BTreeMap::from([("Managed".into(), "org/model".into())]),
4071 )
4072 .unwrap();
4073
4074 assert_eq!(read_receipt(&state, "org/model").unwrap(), Some(receipt));
4075 let _ = crate::retire::revoked_copies(&state, &models);
4076 assert!(quarantine.join("weights").exists(), "nothing resumed");
4077
4078 let _store = ModelManagementStore::new(state, models.clone());
4079 assert!(
4080 std::fs::symlink_metadata(&quarantine).is_err(),
4081 "control: a store resumes it"
4082 );
4083 }
4084
4085 #[cfg(any(target_os = "macos", target_os = "linux"))]
4086 #[test]
4087 fn directory_removal_resumes_after_a_partial_quarantine_delete() {
4088 let (_root, models, store) = supported_platform_store();
4089 let managed = models.join("Managed");
4090 std::fs::create_dir_all(managed.join("sub")).unwrap();
4091 std::fs::write(managed.join("weights"), b"owned").unwrap();
4092 std::fs::write(managed.join("sub").join("more"), b"owned").unwrap();
4093 let receipt = directory_receipt("org/model", managed.clone());
4094 store.write_receipt(&receipt).unwrap();
4095 let quarantine = store.allocate_quarantine("org/model").unwrap();
4096 let journal = RemovalJournal {
4097 model_id: "org/model".into(),
4098 removal_generation: 2,
4099 artifact_identity: artifact_identity(&managed, "org/model").unwrap(),
4100 receipt,
4101 quarantine_path: quarantine.clone(),
4102 };
4103 write_private_json(&store.removal_journal_path("org/model"), &journal).unwrap();
4104 std::fs::rename(&managed, &quarantine).unwrap();
4106 std::fs::remove_file(quarantine.join("weights")).unwrap();
4107
4108 let resumed = store
4109 .begin_mutation("org/model")
4110 .unwrap()
4111 .remove(2)
4112 .unwrap();
4113 assert_eq!(resumed.artifact_kind, ManagedArtifactKind::Directory);
4114 assert!(std::fs::symlink_metadata(&quarantine).is_err());
4115 assert!(no_quarantine_left(&models));
4116
4117 let replayed = store
4122 .begin_mutation("org/model")
4123 .unwrap()
4124 .remove(3)
4125 .unwrap();
4126 assert_eq!(replayed.artifact_kind, ManagedArtifactKind::Directory);
4127 let tombstone = store.load_tombstone("org/model").unwrap().unwrap();
4128 assert_eq!(tombstone.removal_generation, 2);
4129 assert!(replayed.preserved_shared_cache_references.is_empty());
4130 }
4131
4132 #[cfg(any(target_os = "macos", target_os = "linux"))]
4133 #[test]
4134 fn directory_removal_refuses_trees_deeper_than_the_cap() {
4135 let (_root, models, store) = supported_platform_store();
4136 let managed = models.join("Managed");
4137 let mut deep = managed.clone();
4138 for level in 0..(MAX_REMOVAL_DEPTH_FOR_TESTS + 1) {
4139 deep = deep.join(format!("level{level}"));
4140 }
4141 std::fs::create_dir_all(&deep).unwrap();
4142 std::fs::write(deep.join("sentinel"), b"deep").unwrap();
4143 store
4144 .write_receipt(&directory_receipt("org/model", managed.clone()))
4145 .unwrap();
4146
4147 let error = store
4148 .begin_mutation("org/model")
4149 .unwrap()
4150 .remove(2)
4151 .unwrap_err();
4152
4153 let reason = unsafe_managed_reason(&error);
4154 assert!(reason.contains("depth"), "unexpected reason: {reason}");
4155 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4157 let quarantined_leaf = deep
4158 .strip_prefix(&managed)
4159 .map(|rest| journal.quarantine_path.join(rest))
4160 .unwrap();
4161 assert_eq!(
4162 std::fs::read(quarantined_leaf.join("sentinel")).unwrap(),
4163 b"deep"
4164 );
4165 }
4166
4167 #[test]
4171 fn removal_identity_ignores_length_for_directories_only() {
4172 let directory = ArtifactIdentity {
4173 kind: ManagedArtifactKind::Directory,
4174 device: 7,
4175 inode: 42,
4176 file_len: 96,
4177 };
4178 let shrunk = ArtifactIdentity {
4179 file_len: 32,
4180 ..directory.clone()
4181 };
4182 assert!(directory.matches_for_removal(&shrunk));
4183 assert!(!directory.matches_for_removal(&ArtifactIdentity {
4184 inode: 43,
4185 ..directory.clone()
4186 }));
4187 assert!(!directory.matches_for_removal(&ArtifactIdentity {
4188 device: 8,
4189 ..directory.clone()
4190 }));
4191
4192 let file = ArtifactIdentity {
4193 kind: ManagedArtifactKind::File,
4194 device: 7,
4195 inode: 42,
4196 file_len: 96,
4197 };
4198 assert!(file.matches_for_removal(&file));
4199 assert!(!file.matches_for_removal(&ArtifactIdentity {
4200 file_len: 95,
4201 ..file.clone()
4202 }));
4203 assert!(!file.matches_for_removal(&ArtifactIdentity {
4204 kind: ManagedArtifactKind::Symlink,
4205 ..file.clone()
4206 }));
4207 }
4208
4209 #[cfg(any(target_os = "macos", target_os = "linux"))]
4210 #[test]
4211 fn same_filesystem_predicate_compares_devices_only() {
4212 use super::object_bound::{same_filesystem, Identity};
4213 let root = Identity {
4214 device: 5,
4215 inode: 1,
4216 };
4217 assert!(same_filesystem(
4218 &root,
4219 &Identity {
4220 device: 5,
4221 inode: 999
4222 }
4223 ));
4224 assert!(!same_filesystem(
4225 &root,
4226 &Identity {
4227 device: 6,
4228 inode: 1
4229 }
4230 ));
4231 }
4232
4233 #[cfg(any(target_os = "macos", target_os = "linux"))]
4234 #[test]
4235 fn directory_removal_refuses_entries_added_after_capture() {
4236 let (_root, models, store) = supported_platform_store();
4237 let managed = models.join("Managed");
4238 std::fs::create_dir_all(managed.join("sub")).unwrap();
4239 std::fs::write(managed.join("weights"), b"owned").unwrap();
4240 std::fs::write(managed.join("sub").join("more"), b"owned").unwrap();
4241 let store = store.with_removal_hook(|phase, quarantine| {
4242 if phase == RemovalPhase::AfterCapture {
4243 std::fs::write(quarantine.join("sub").join("late"), b"not captured").unwrap();
4244 }
4245 });
4246 store
4247 .write_receipt(&directory_receipt("org/model", managed.clone()))
4248 .unwrap();
4249
4250 let error = store
4251 .begin_mutation("org/model")
4252 .unwrap()
4253 .remove(2)
4254 .unwrap_err();
4255
4256 let reason = unsafe_managed_reason(&error);
4257 assert!(
4258 reason.contains("not captured"),
4259 "unexpected reason: {reason}"
4260 );
4261 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4262 let late = journal.quarantine_path.join("sub").join("late");
4263 assert_eq!(std::fs::read(&late).unwrap(), b"not captured");
4264
4265 std::fs::remove_file(&late).unwrap();
4267 let store = store.with_removal_hook(|_, _| {});
4268 let resumed = store
4269 .begin_mutation("org/model")
4270 .unwrap()
4271 .remove(2)
4272 .unwrap();
4273 assert_eq!(resumed.artifact_kind, ManagedArtifactKind::Directory);
4274 assert!(no_quarantine_left(&models));
4275 }
4276
4277 #[cfg(any(target_os = "macos", target_os = "linux"))]
4278 #[test]
4279 fn directory_removal_refuses_root_swap_after_rename() {
4280 let (_root, models, store) = supported_platform_store();
4281 let managed = models.join("Managed");
4282 std::fs::create_dir_all(&managed).unwrap();
4283 std::fs::write(managed.join("weights"), b"owned").unwrap();
4284 let displaced = models.join("displaced");
4285 let store = store.with_removal_hook(move |phase, quarantine| {
4286 if phase == RemovalPhase::AfterRename {
4287 std::fs::rename(quarantine, &displaced).unwrap();
4288 std::fs::create_dir_all(quarantine).unwrap();
4289 std::fs::write(quarantine.join("sentinel"), b"replacement").unwrap();
4290 }
4291 });
4292 store
4293 .write_receipt(&directory_receipt("org/model", managed.clone()))
4294 .unwrap();
4295
4296 let error = store
4297 .begin_mutation("org/model")
4298 .unwrap()
4299 .remove(2)
4300 .unwrap_err();
4301
4302 let reason = unsafe_managed_reason(&error);
4303 assert!(
4304 reason.contains("identity changed"),
4305 "unexpected reason: {reason}"
4306 );
4307 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4308 assert_eq!(
4309 std::fs::read(journal.quarantine_path.join("sentinel")).unwrap(),
4310 b"replacement"
4311 );
4312 assert_eq!(
4313 std::fs::read(models.join("displaced").join("weights")).unwrap(),
4314 b"owned"
4315 );
4316 }
4317
4318 #[cfg(any(target_os = "macos", target_os = "linux"))]
4319 #[test]
4320 fn directory_symlink_planted_after_rename_is_unlinked_not_followed() {
4321 use std::os::unix::fs::symlink;
4322
4323 let (root, models, store) = supported_platform_store();
4324 let outside = root.path().join("outside");
4325 std::fs::create_dir_all(&outside).unwrap();
4326 std::fs::write(outside.join("sentinel"), b"outside").unwrap();
4327 let managed = models.join("Managed");
4328 std::fs::create_dir_all(&managed).unwrap();
4329 std::fs::write(managed.join("weights"), b"owned").unwrap();
4330 let target = outside.clone();
4331 let store = store.with_removal_hook(move |phase, quarantine| {
4332 if phase == RemovalPhase::AfterRename {
4333 symlink(&target, quarantine.join("escape")).unwrap();
4334 }
4335 });
4336 store
4337 .write_receipt(&directory_receipt("org/model", managed.clone()))
4338 .unwrap();
4339
4340 let result = store.begin_mutation("org/model").unwrap().remove(2);
4341
4342 assert_eq!(
4343 std::fs::read(outside.join("sentinel")).map_err(|e| e.to_string()),
4344 Ok(b"outside".to_vec()),
4345 "outside sentinel was deleted; removal result: {result:?}"
4346 );
4347 assert!(result.is_ok(), "{result:?}");
4348 assert!(std::fs::symlink_metadata(&managed).is_err());
4349 assert!(no_quarantine_left(&models));
4350 }
4351
4352 #[cfg(any(target_os = "macos", target_os = "linux"))]
4353 #[test]
4354 fn directory_removal_refuses_root_swap_between_path_check_and_descriptor_open() {
4355 let (_root, models, store) = supported_platform_store();
4356 let managed = models.join("Managed");
4357 std::fs::create_dir_all(&managed).unwrap();
4358 std::fs::write(managed.join("weights"), b"owned").unwrap();
4359 let displaced = models.join("displaced");
4360 let store = store.with_removal_hook(move |phase, quarantine| {
4361 if phase == RemovalPhase::BeforeRootOpen {
4362 std::fs::rename(quarantine, &displaced).unwrap();
4363 std::fs::create_dir_all(quarantine).unwrap();
4364 std::fs::write(quarantine.join("sentinel"), b"replacement").unwrap();
4365 }
4366 });
4367 store
4368 .write_receipt(&directory_receipt("org/model", managed.clone()))
4369 .unwrap();
4370
4371 let error = store
4372 .begin_mutation("org/model")
4373 .unwrap()
4374 .remove(2)
4375 .unwrap_err();
4376
4377 let reason = unsafe_managed_reason(&error);
4378 assert_eq!(
4379 reason,
4380 "quarantined directory identity changed before deletion"
4381 );
4382 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4383 assert_eq!(
4384 std::fs::read(journal.quarantine_path.join("sentinel")).unwrap(),
4385 b"replacement"
4386 );
4387 assert_eq!(
4388 std::fs::read(models.join("displaced").join("weights")).unwrap(),
4389 b"owned"
4390 );
4391 }
4392
4393 #[cfg(any(target_os = "macos", target_os = "linux"))]
4394 #[test]
4395 fn directory_removal_refuses_when_the_quarantine_vanishes_after_our_own_rename() {
4396 let (_root, models, store) = supported_platform_store();
4397 let managed = models.join("Managed");
4398 std::fs::create_dir_all(&managed).unwrap();
4399 std::fs::write(managed.join("weights"), b"owned").unwrap();
4400 let managed_for_hook = managed.clone();
4401 let store = store.with_removal_hook(move |phase, quarantine| {
4402 if phase == RemovalPhase::AfterRename {
4403 std::fs::rename(quarantine, &managed_for_hook).unwrap();
4405 }
4406 });
4407 store
4408 .write_receipt(&directory_receipt("org/model", managed.clone()))
4409 .unwrap();
4410
4411 let error = store
4412 .begin_mutation("org/model")
4413 .unwrap()
4414 .remove(2)
4415 .unwrap_err();
4416
4417 let reason = unsafe_managed_reason(&error);
4418 assert!(reason.contains("vanished"), "unexpected reason: {reason}");
4419 assert_eq!(std::fs::read(managed.join("weights")).unwrap(), b"owned");
4420 assert!(store.load_receipt("org/model").unwrap().is_some());
4421 assert!(store.load_removal_journal("org/model").unwrap().is_some());
4422 }
4423
4424 #[cfg(any(target_os = "macos", target_os = "linux"))]
4425 #[test]
4426 fn first_attempt_vanish_refusal_names_the_retry_that_completes_the_removal() {
4427 let (_root, models, store) = supported_platform_store();
4428 let managed = models.join("Managed");
4429 std::fs::create_dir_all(&managed).unwrap();
4430 std::fs::write(managed.join("weights"), b"owned").unwrap();
4431 let managed_for_hook = managed.clone();
4432 let sabotaged = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
4433 let store = store.with_removal_hook(move |phase, quarantine| {
4434 if phase == RemovalPhase::AfterRename
4438 && !sabotaged.swap(true, std::sync::atomic::Ordering::SeqCst)
4439 {
4440 std::fs::rename(quarantine, &managed_for_hook).unwrap();
4441 }
4442 });
4443 store
4444 .write_receipt(&directory_receipt("org/model", managed.clone()))
4445 .unwrap();
4446
4447 let error = store
4448 .begin_mutation("org/model")
4449 .unwrap()
4450 .remove(2)
4451 .unwrap_err();
4452
4453 assert_eq!(
4457 unsafe_managed_reason(&error),
4458 "quarantine vanished before deletion; the removal journal is retained, so retrying the removal resumes and completes it"
4459 );
4460 assert!(store.load_removal_journal("org/model").unwrap().is_some());
4461
4462 let retried = store
4464 .begin_mutation("org/model")
4465 .unwrap()
4466 .remove(2)
4467 .unwrap();
4468
4469 assert_eq!(retried.artifact_kind, ManagedArtifactKind::Directory);
4470 assert!(std::fs::symlink_metadata(&managed).is_err());
4471 assert!(no_quarantine_left(&models));
4472 assert!(store.load_removal_journal("org/model").unwrap().is_none());
4473 assert!(store.load_receipt("org/model").unwrap().is_none());
4474 }
4475
4476 #[cfg(any(target_os = "macos", target_os = "linux"))]
4477 #[test]
4478 fn directory_removal_refuses_a_captured_file_replaced_before_unlink() {
4479 let (root, models, store) = supported_platform_store();
4480 let outside = root.path().join("outside");
4481 std::fs::create_dir_all(&outside).unwrap();
4482 let managed = models.join("Managed");
4483 std::fs::create_dir_all(&managed).unwrap();
4484 std::fs::write(managed.join("weights"), b"owned").unwrap();
4485 let parked = outside.join("parked");
4486 let store = store.with_removal_hook(move |phase, quarantine| {
4487 if phase == RemovalPhase::AfterCapture {
4488 std::fs::rename(quarantine.join("weights"), &parked).unwrap();
4489 std::fs::write(quarantine.join("weights"), b"victim").unwrap();
4490 }
4491 });
4492 store
4493 .write_receipt(&directory_receipt("org/model", managed.clone()))
4494 .unwrap();
4495
4496 let error = store
4497 .begin_mutation("org/model")
4498 .unwrap()
4499 .remove(2)
4500 .unwrap_err();
4501
4502 let reason = unsafe_managed_reason(&error);
4503 assert!(
4504 reason.contains("identity changed"),
4505 "unexpected reason: {reason}"
4506 );
4507 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4508 assert_eq!(
4509 std::fs::read(journal.quarantine_path.join("weights")).unwrap(),
4510 b"victim"
4511 );
4512 assert_eq!(std::fs::read(outside.join("parked")).unwrap(), b"owned");
4513 }
4514
4515 #[cfg(any(target_os = "macos", target_os = "linux"))]
4516 #[test]
4517 fn directory_removal_refuses_a_captured_directory_replaced_before_rmdir() {
4518 let (root, models, store) = supported_platform_store();
4519 let outside = root.path().join("outside");
4520 std::fs::create_dir_all(&outside).unwrap();
4521 let managed = models.join("Managed");
4522 std::fs::create_dir_all(managed.join("sub")).unwrap();
4523 std::fs::write(managed.join("sub").join("more"), b"owned").unwrap();
4524 let parked = outside.join("parked-sub");
4525 let parked_for_hook = parked.clone();
4526 let store = store.with_removal_hook(move |phase, quarantine| {
4527 if phase == RemovalPhase::AfterCapture {
4528 std::fs::rename(quarantine.join("sub"), &parked_for_hook).unwrap();
4529 std::fs::create_dir_all(quarantine.join("sub")).unwrap();
4530 std::fs::write(quarantine.join("sub").join("victim"), b"victim").unwrap();
4531 }
4532 });
4533 store
4534 .write_receipt(&directory_receipt("org/model", managed.clone()))
4535 .unwrap();
4536
4537 let error = store
4538 .begin_mutation("org/model")
4539 .unwrap()
4540 .remove(2)
4541 .unwrap_err();
4542
4543 let reason = unsafe_managed_reason(&error);
4544 assert!(
4545 reason.contains("identity changed"),
4546 "unexpected reason: {reason}"
4547 );
4548 let journal = store.load_removal_journal("org/model").unwrap().unwrap();
4549 assert_eq!(
4550 std::fs::read(journal.quarantine_path.join("sub").join("victim")).unwrap(),
4551 b"victim"
4552 );
4553 assert!(std::fs::symlink_metadata(&parked).is_ok());
4556 }
4557
4558 #[cfg(any(target_os = "macos", target_os = "linux"))]
4559 #[test]
4560 fn directory_removal_accepts_a_tree_at_the_depth_cap() {
4561 let (_root, models, store) = supported_platform_store();
4562 let managed = models.join("Managed");
4563 let mut deep = managed.clone();
4564 for level in 0..MAX_REMOVAL_DEPTH_FOR_TESTS {
4565 deep = deep.join(format!("level{level}"));
4566 }
4567 std::fs::create_dir_all(&deep).unwrap();
4568 std::fs::write(deep.join("sentinel"), b"deep").unwrap();
4569 store
4570 .write_receipt(&directory_receipt("org/model", managed.clone()))
4571 .unwrap();
4572
4573 let result = store
4574 .begin_mutation("org/model")
4575 .unwrap()
4576 .remove(2)
4577 .unwrap();
4578
4579 assert_eq!(result.artifact_kind, ManagedArtifactKind::Directory);
4580 assert!(std::fs::symlink_metadata(&managed).is_err());
4581 assert!(no_quarantine_left(&models));
4582 }
4583
4584 #[cfg(any(target_os = "macos", target_os = "linux"))]
4585 #[test]
4586 fn directory_removal_refuses_when_the_quarantine_vanishes_before_the_descriptor_open() {
4587 let (_root, models, store) = supported_platform_store();
4588 let managed = models.join("Managed");
4589 std::fs::create_dir_all(&managed).unwrap();
4590 std::fs::write(managed.join("weights"), b"owned").unwrap();
4591 let managed_for_hook = managed.clone();
4592 let store = store.with_removal_hook(move |phase, quarantine| {
4593 if phase == RemovalPhase::BeforeRootOpen {
4594 std::fs::rename(quarantine, &managed_for_hook).unwrap();
4597 }
4598 });
4599 store
4600 .write_receipt(&directory_receipt("org/model", managed.clone()))
4601 .unwrap();
4602
4603 let error = store
4604 .begin_mutation("org/model")
4605 .unwrap()
4606 .remove(2)
4607 .unwrap_err();
4608
4609 assert_eq!(
4614 unsafe_managed_reason(&error),
4615 "quarantine vanished before deletion; the removal journal is retained, so retrying the removal resumes and completes it"
4616 );
4617 assert_eq!(std::fs::read(managed.join("weights")).unwrap(), b"owned");
4618 assert!(store.load_receipt("org/model").unwrap().is_some());
4619 assert!(store.load_removal_journal("org/model").unwrap().is_some());
4620 }
4621
4622 #[cfg(any(target_os = "macos", target_os = "linux"))]
4629 #[test]
4630 fn retry_after_the_quarantine_was_parked_elsewhere_settles_without_deleting_it() {
4631 let (root, models, store) = supported_platform_store();
4632 let outside = root.path().join("outside");
4633 std::fs::create_dir_all(&outside).unwrap();
4634 let managed = models.join("Managed");
4635 std::fs::create_dir_all(&managed).unwrap();
4636 std::fs::write(managed.join("weights"), b"owned").unwrap();
4637 let parked = outside.join("parked");
4638 let parked_for_hook = parked.clone();
4639 let store = store.with_removal_hook(move |phase, quarantine| {
4640 if phase == RemovalPhase::AfterRename {
4641 std::fs::rename(quarantine, &parked_for_hook).unwrap();
4643 }
4644 });
4645 store
4646 .write_receipt(&directory_receipt("org/model", managed.clone()))
4647 .unwrap();
4648
4649 let error = store
4650 .begin_mutation("org/model")
4651 .unwrap()
4652 .remove(2)
4653 .unwrap_err();
4654
4655 assert_eq!(
4656 unsafe_managed_reason(&error),
4657 "quarantine vanished before deletion; the removal journal is retained, so retrying the removal resumes and completes it"
4658 );
4659 assert!(store.load_removal_journal("org/model").unwrap().is_some());
4660
4661 let retried = store
4662 .begin_mutation("org/model")
4663 .unwrap()
4664 .remove(2)
4665 .unwrap();
4666
4667 assert_eq!(retried.artifact_kind, ManagedArtifactKind::Directory);
4668 assert!(store.load_removal_journal("org/model").unwrap().is_none());
4669 assert!(store.load_receipt("org/model").unwrap().is_none());
4670 assert!(no_quarantine_left(&models));
4671 assert_eq!(std::fs::read(parked.join("weights")).unwrap(), b"owned");
4673 }
4674}