diff --git a/Cargo.lock b/Cargo.lock index daa43d2..e5de3cf 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -38,7 +38,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -49,7 +49,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -73,6 +73,21 @@ version = "0.22.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" +[[package]] +name = "bitflags" +version = "2.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2261d10cca569e4643e526d8dc2e62e433cc8aba21ab764233731f8d369bf394" + +[[package]] +name = "block-buffer" +version = "0.10.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71" +dependencies = [ + "generic-array", +] + [[package]] name = "bstr" version = "1.13.0" @@ -84,6 +99,12 @@ dependencies = [ "serde_core", ] +[[package]] +name = "cfg-if" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2fd1289c04a9ea8cb22300a459a72a385d7c73d3259e2ed7dcb2af674838cfa9" + [[package]] name = "clap" version = "4.6.4" @@ -130,18 +151,85 @@ version = "1.0.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" +[[package]] +name = "cpufeatures" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ed5838eebb26a2bb2e58f6d5b5316989ae9d08bab10e0e6d103e656d1b0280" +dependencies = [ + "libc", +] + +[[package]] +name = "crypto-common" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bfb12502f3fc46cca1bb51ac28df9d618d813cdc3d2f25b9fe775a34af26bb3" +dependencies = [ + "generic-array", + "typenum", +] + [[package]] name = "difflib" version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6184e33543162437515c2e2b48714794e37845ec9851711914eec9d308f6ebe8" +[[package]] +name = "digest" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" +dependencies = [ + "block-buffer", + "crypto-common", +] + [[package]] name = "dyn-clone" version = "1.0.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555" +[[package]] +name = "errno" +version = "0.3.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "778e2ac28f6c47af28e4907f13ffd1e1ddbd400980a9abd7c8df189bf578a5ad" +dependencies = [ + "libc", + "windows-sys 0.60.2", +] + +[[package]] +name = "fastrand" +version = "2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" + +[[package]] +name = "generic-array" +version = "0.14.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a" +dependencies = [ + "typenum", + "version_check", +] + +[[package]] +name = "getrandom" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26145e563e54f2cadc477553f1ec5ee650b00862f0a58bcd12cbdc5f0ea2d2f4" +dependencies = [ + "cfg-if", + "libc", + "r-efi", + "wasi", +] + [[package]] name = "ghidra-cli" version = "0.1.0" @@ -149,9 +237,14 @@ dependencies = [ "assert_cmd", "base64", "clap", + "getrandom", + "hex", + "rustix", "schemars", "serde", "serde_json", + "sha2", + "tempfile", "thiserror", ] @@ -161,6 +254,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hex" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" + [[package]] name = "is_terminal_polyfill" version = "1.70.2" @@ -179,12 +278,24 @@ version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "linux-raw-sys" +version = "0.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cd945864f07fe9f5371a27ad7b52a172b4b499999f1d97574c9fa68373937e12" + [[package]] name = "memchr" version = "2.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cf8baf1c55e62ffcace7a9f06f4bd9cd3f0c4beb022d3b367256b91b87513d98" +[[package]] +name = "once_cell" +version = "1.21.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d" + [[package]] name = "once_cell_polyfill" version = "1.70.2" @@ -236,6 +347,12 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" + [[package]] name = "ref-cast" version = "1.0.26" @@ -262,6 +379,19 @@ version = "0.4.16" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8fcfdb36bda0c880c5931cdc7a2bcdc8ba4556847b9d912bca70bc94708711ad" +[[package]] +name = "rustix" +version = "1.0.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "11181fbabf243db407ef8df94a6ce0b2f9a733bd8be4ad02b4eda9602296cac8" +dependencies = [ + "bitflags", + "errno", + "libc", + "linux-raw-sys", + "windows-sys 0.60.2", +] + [[package]] name = "schemars" version = "1.0.4" @@ -341,6 +471,17 @@ dependencies = [ "zmij", ] +[[package]] +name = "sha2" +version = "0.10.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "strsim" version = "0.11.1" @@ -369,6 +510,19 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "tempfile" +version = "3.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15b61f8f20e3a6f7e0649d825294eaf317edce30f82cf6026e7e4cb9222a7d1e" +dependencies = [ + "fastrand", + "getrandom", + "once_cell", + "rustix", + "windows-sys 0.60.2", +] + [[package]] name = "termtree" version = "0.5.1" @@ -395,6 +549,12 @@ dependencies = [ "syn 3.0.3", ] +[[package]] +name = "typenum" +version = "1.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1dccffe3ce07af9386bfd29e80c0ab1a8205a2fc34e4bcd40364df902cfa8f3f" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -407,6 +567,12 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "version_check" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" + [[package]] name = "wait-timeout" version = "0.2.1" @@ -416,21 +582,116 @@ dependencies = [ "libc", ] +[[package]] +name = "wasi" +version = "0.14.3+wasi-0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6a51ae83037bdd272a9e28ce236db8c07016dd0d50c27038b3f407533c030c95" +dependencies = [ + "wit-bindgen", +] + +[[package]] +name = "windows-link" +version = "0.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e6ad25900d524eaabdbbb96d20b4311e1e7ae1699af4fb28c17ae66c80d798a" + [[package]] name = "windows-link" version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" +[[package]] +name = "windows-sys" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" +dependencies = [ + "windows-targets", +] + [[package]] name = "windows-sys" version = "0.61.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" dependencies = [ - "windows-link", + "windows-link 0.2.1", ] +[[package]] +name = "windows-targets" +version = "0.53.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d5fe6031c4041849d7c496a8ded650796e7b6ecc19df1a431c1a363342e5dc91" +dependencies = [ + "windows-link 0.1.3", + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_gnullvm", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86b8d5f90ddd19cb4a147a5fa63ca848db3df085e25fee3cc10b39b6eebae764" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7651a1f62a11b8cbd5e0d42526e55f2c99886c77e007179efff86c2b137e66c" + +[[package]] +name = "windows_i686_gnu" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1dc67659d35f387f5f6c479dc4e28f1d4bb90ddd1a5d3da2e5d97b42d6272c3" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ce6ccbdedbf6d6354471319e781c0dfef054c81fbc7cf83f338a4296c0cae11" + +[[package]] +name = "windows_i686_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "581fee95406bb13382d2f65cd4a908ca7b1e4c2f1917f143ba16efe98a589b5d" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e55b5ac9ea33f2fc1716d1742db15574fd6fc8dadc51caab1c16a3d3b4190ba" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a6e035dd0599267ce1ee132e51c27dd29437f63325753051e71dd9e42406c57" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "271414315aff87387382ec3d271b52d7ae78726f5d44ac98b4f4030c91880486" + +[[package]] +name = "wit-bindgen" +version = "0.45.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "052283831dbae3d879dc7f51f3d92703a316ca49f91540417d38591826127814" + [[package]] name = "zmij" version = "1.0.23" diff --git a/Cargo.toml b/Cargo.toml index efabe06..a34eae3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,13 +18,18 @@ path = "src/main.rs" [dependencies] base64 = "0.22.1" clap = { version = "4.5.60", features = ["derive"] } +getrandom = "0.3.3" +hex = "0.4.3" +rustix = { version = "1.0.8", features = ["fs"] } schemars = { version = "=1.0.4", features = ["derive"] } serde = { version = "1.0.228", features = ["derive"] } serde_json = "1.0.149" +sha2 = "0.10.9" thiserror = "2.0.18" [dev-dependencies] assert_cmd = "2.1.2" +tempfile = "3.21.0" [lints.rust] unsafe_code = "forbid" diff --git a/src/lib.rs b/src/lib.rs index 858c23b..d5e838a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -8,5 +8,6 @@ pub mod operation; pub mod output; pub mod protocol; pub mod schema; +pub mod store; pub use error::{AppError, ErrorCode, ExitStatus}; diff --git a/src/store/analysis.rs b/src/store/analysis.rs new file mode 100644 index 0000000..732f5d8 --- /dev/null +++ b/src/store/analysis.rs @@ -0,0 +1,553 @@ +use std::{ + collections::BTreeMap, + fs::{self, File}, + io::Read as _, + path::{Component, Path, PathBuf}, + str::FromStr as _, +}; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use sha2::{Digest as _, Sha256}; + +use crate::{domain::Digest, operation::AnalysisProfile}; + +use super::{ + AnalysisLock, AnalysisStore, StoreError, StoreErrorKind, + filesystem::{create_private_directory, random_id, sync_directory}, +}; + +/// Complete immutable Analysis manifest. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct AnalysisManifest { + /// Manifest format version, exactly 1. + pub manifest_version: u32, + /// Completion marker, exactly `complete`. + pub state: String, + /// Staged Sample identity and byte length. + pub sample: ManifestSample, + /// Analysis Profile identity and compatibility facts. + pub analysis_profile: AnalysisProfileManifest, + /// Tool and invocation that completed the Analysis. + pub created_by: CreatedBy, + /// Fixed project name, path, and complete file inventory. + pub project: AnalysisProject, +} + +/// Sample facts recorded in an Analysis manifest. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct ManifestSample { + /// SHA-256 Sample identity. + pub sha256: Digest, + /// Exact Sample byte size. + pub size_bytes: u64, +} + +/// Compatibility facts recorded alongside the Profile digest. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct AnalysisProfileManifest { + /// Canonical Analysis Profile digest. + pub sha256: Digest, + /// Exact Ghidra version. + pub ghidra_version: String, + /// Rust/Java adapter protocol version. + pub adapter_protocol_version: u32, + /// Fully resolved Target Specification facts. + pub target: BTreeMap, + /// Digest of the fully resolved analyzer-option document. + pub analyzer_options_sha256: Digest, +} + +/// Creator provenance for one immutable Analysis. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct CreatedBy { + /// CLI version. + pub ghidr_version: String, + /// Canonical random invocation ID. + pub invocation_id: String, + /// RFC3339 UTC completion timestamp supplied by orchestration. + pub completed_at: String, +} + +/// Fixed project metadata and complete inventory. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct AnalysisProject { + /// Fixed project basename, exactly `analysis`. + pub name: String, + /// Fixed relative project directory, exactly `project`. + pub path: String, + /// Every regular file in bytewise relative-path order. + pub files: Vec, +} + +/// One exact project file inventory entry. +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct InventoryFile { + /// UTF-8 relative path beneath `project`. + pub path: String, + /// Exact file length. + pub size_bytes: u64, + /// SHA-256 of exact file bytes. + pub sha256: Digest, +} + +/// Outcome of moving corrupt active Analysis data out of visibility. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct QuarantineOutcome { + /// Unique retained quarantine directory. + pub path: PathBuf, + /// Required next decision after successful quarantine. + pub decision: RebuildDecision, +} + +/// Single-rebuild state machine primitive. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RebuildDecision { + /// One fresh build is permitted. + RebuildOnce, + /// Recovery was already attempted; surface corruption. + FailCorrupt, +} + +/// Returns the only permitted recovery action for the invocation state. +#[must_use] +pub const fn rebuild_decision(rebuild_already_attempted: bool) -> RebuildDecision { + if rebuild_already_attempted { + RebuildDecision::FailCorrupt + } else { + RebuildDecision::RebuildOnce + } +} + +/// Produces the canonical Profile digest after enforcing sorted unique options. +pub fn analysis_profile_digest(profile: &AnalysisProfile) -> Result { + if profile.profile_version != 1 { + return Err(StoreError::new( + StoreErrorKind::CorruptStore, + "Analysis Profile version must be exactly 1", + )); + } + validate_options(&profile.loader_options, "loader")?; + validate_options(&profile.analyzer_options, "analyzer")?; + let canonical = serde_json::to_value(profile) + .and_then(|value| serde_json::to_vec(&value)) + .map_err(|error| { + StoreError::new( + StoreErrorKind::CorruptStore, + format!("Profile serialization failed: {error}"), + ) + })?; + let digest = Sha256::digest(&canonical); + digest_bytes(&digest) +} + +/// Inventories a complete project tree in manifest bytewise path order. +pub fn inventory_project(project: &Path) -> Result, StoreError> { + reject_tree_root(project)?; + inventory_tree(project, project) +} + +/// Atomically writes the completed manifest into a private staged Analysis directory. +pub fn write_analysis_manifest( + staged_analysis: &Path, + manifest: &AnalysisManifest, +) -> Result<(), StoreError> { + super::filesystem::atomic_json_write( + staged_analysis, + &staged_analysis.join("manifest.json"), + manifest, + ) +} + +fn validate_options( + options: &[crate::operation::ProfileOption], + kind: &str, +) -> Result<(), StoreError> { + if options + .windows(2) + .any(|pair| pair[0].name.as_bytes() >= pair[1].name.as_bytes()) + { + return Err(StoreError::new( + StoreErrorKind::CorruptStore, + format!("{kind} options must be bytewise sorted with unique names"), + )); + } + Ok(()) +} + +/// Validates manifest identity and every regular project file without following symlinks. +pub fn validate_analysis( + analysis_directory: &Path, + expected_sample: &Digest, + expected_profile: &Digest, +) -> Result { + reject_tree_root(analysis_directory)?; + let manifest_path = analysis_directory.join("manifest.json"); + require_regular(&manifest_path)?; + let bytes = fs::read(&manifest_path).map_err(|error| StoreError::io(&manifest_path, error))?; + let manifest: AnalysisManifest = serde_json::from_slice(&bytes).map_err(|error| { + StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + format!("invalid Analysis manifest: {error}"), + ) + })?; + if manifest.manifest_version != 1 + || manifest.state != "complete" + || &manifest.sample.sha256 != expected_sample + || &manifest.analysis_profile.sha256 != expected_profile + || manifest.project.name != "analysis" + || manifest.project.path != "project" + { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + "Analysis manifest identity or fixed fields do not match its location", + )); + } + validate_inventory(&analysis_directory.join("project"), &manifest.project.files)?; + let mut root_entries = read_sorted(analysis_directory)?; + let root_names: Vec<_> = root_entries + .drain(..) + .map(|entry| entry.file_name()) + .collect(); + if root_names != ["manifest.json", "project"].map(std::ffi::OsString::from) { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + analysis_directory, + "Analysis directory contains missing or unexpected entries", + )); + } + Ok(manifest) +} + +impl AnalysisStore { + /// Validates and atomically promotes a staged complete Analysis. + pub fn promote_analysis( + &self, + staged_analysis: &Path, + sample: &Digest, + profile: &Digest, + lock: &AnalysisLock, + ) -> Result { + if !lock.guards(sample, profile) { + return Err(StoreError::new( + StoreErrorKind::AnalysisBusy, + "lock token does not guard the promoted Analysis identity", + )); + } + validate_analysis(staged_analysis, sample, profile)?; + let sample_directory = self.path("analyses").join(sample.as_str()); + create_private_directory(&sample_directory)?; + let destination = sample_directory.join(profile.as_str()); + if fs::symlink_metadata(&destination).is_ok() { + return Err(StoreError::at( + StoreErrorKind::ImmutableConflict, + &destination, + "completed Analysis destination already exists", + )); + } + fs::rename(staged_analysis, &destination) + .map_err(|error| StoreError::io(&destination, error))?; + sync_directory(&sample_directory)?; + Ok(destination) + } + + /// Atomically removes corrupt Analysis data from active visibility. + pub fn quarantine_analysis( + &self, + sample: &Digest, + profile: &Digest, + lock: &AnalysisLock, + rebuild_already_attempted: bool, + ) -> Result { + if !lock.guards(sample, profile) { + return Err(StoreError::new( + StoreErrorKind::AnalysisBusy, + "lock token does not guard the quarantined Analysis identity", + )); + } + let active = self + .path("analyses") + .join(sample.as_str()) + .join(profile.as_str()); + let parent = self + .path("quarantine") + .join(sample.as_str()) + .join(profile.as_str()); + create_private_directory(&parent)?; + let destination = parent.join(random_id()?); + fs::rename(&active, &destination).map_err(|error| StoreError::io(&active, error))?; + sync_directory(&parent)?; + Ok(QuarantineOutcome { + path: destination, + decision: rebuild_decision(rebuild_already_attempted), + }) + } +} + +fn validate_inventory(project: &Path, recorded: &[InventoryFile]) -> Result<(), StoreError> { + reject_tree_root(project)?; + if recorded + .windows(2) + .any(|pair| pair[0].path.as_bytes() >= pair[1].path.as_bytes()) + { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + project, + "project inventory is not bytewise sorted and unique", + )); + } + let actual = inventory_tree(project, project)?; + if actual != recorded { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + project, + "project inventory size, digest, or membership mismatch", + )); + } + Ok(()) +} + +fn inventory_tree(root: &Path, directory: &Path) -> Result, StoreError> { + let mut inventory = Vec::new(); + for entry in read_sorted(directory)? { + let path = entry.path(); + let metadata = fs::symlink_metadata(&path).map_err(|error| StoreError::io(&path, error))?; + if metadata.file_type().is_symlink() { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &path, + "symlink in Analysis project", + )); + } + if metadata.is_dir() { + inventory.extend(inventory_tree(root, &path)?); + } else if metadata.is_file() { + let relative = path.strip_prefix(root).map_err(|_| { + StoreError::at( + StoreErrorKind::CorruptStore, + &path, + "project path escaped inventory root", + ) + })?; + if relative + .components() + .any(|component| !matches!(component, Component::Normal(_))) + { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &path, + "invalid inventory path", + )); + } + let Some(relative) = relative.to_str() else { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &path, + "non-UTF-8 project filename", + )); + }; + inventory.push(InventoryFile { + path: relative.to_owned(), + size_bytes: metadata.len(), + sha256: hash_file(&path)?, + }); + } else { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &path, + "special file in Analysis project", + )); + } + } + inventory.sort_by(|left, right| left.path.as_bytes().cmp(right.path.as_bytes())); + Ok(inventory) +} + +pub(crate) fn hash_file(path: &Path) -> Result { + let mut file = File::open(path).map_err(|error| StoreError::io(path, error))?; + let mut hasher = Sha256::new(); + let mut buffer = vec![0_u8; 1024 * 1024]; + loop { + let count = file + .read(&mut buffer) + .map_err(|error| StoreError::io(path, error))?; + if count == 0 { + break; + } + hasher.update(&buffer[..count]); + } + digest_bytes(&hasher.finalize()) +} + +pub(crate) fn digest_bytes(bytes: &[u8]) -> Result { + Digest::from_str(&hex::encode(bytes)).map_err(|error| { + StoreError::new( + StoreErrorKind::Io, + format!("computed digest was invalid: {error}"), + ) + }) +} + +fn reject_tree_root(path: &Path) -> Result<(), StoreError> { + let metadata = fs::symlink_metadata(path).map_err(|error| StoreError::io(path, error))?; + if !metadata.is_dir() || metadata.file_type().is_symlink() { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + path, + "expected ordinary directory", + )); + } + Ok(()) +} + +fn require_regular(path: &Path) -> Result<(), StoreError> { + let metadata = fs::symlink_metadata(path).map_err(|error| StoreError::io(path, error))?; + if !metadata.is_file() || metadata.file_type().is_symlink() { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + path, + "expected ordinary file", + )); + } + Ok(()) +} + +fn read_sorted(path: &Path) -> Result, StoreError> { + let mut entries = fs::read_dir(path) + .map_err(|error| StoreError::io(path, error))? + .collect::, _>>() + .map_err(|error| StoreError::io(path, error))?; + entries.sort_by_key(fs::DirEntry::file_name); + Ok(entries) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{collections::BTreeMap, fs, str::FromStr as _}; + + use tempfile::TempDir; + + use super::{ + AnalysisManifest, AnalysisProfileManifest, AnalysisProject, CreatedBy, InventoryFile, + ManifestSample, RebuildDecision, hash_file, rebuild_decision, validate_analysis, + }; + use crate::{ + domain::Digest, + store::{AnalysisLock, AnalysisStore, StoreEnvironment, resolve_store}, + }; + + fn digest(byte: char) -> Digest { + Digest::from_str(&byte.to_string().repeat(64)).expect("digest") + } + + fn write_analysis(root: &Path, sample: &Digest, profile: &Digest) { + let project = root.join("project"); + fs::create_dir_all(&project).expect("project"); + let project_file = project.join("analysis.gpr"); + fs::write(&project_file, b"project").expect("project file"); + let manifest = AnalysisManifest { + manifest_version: 1, + state: "complete".to_owned(), + sample: ManifestSample { + sha256: sample.clone(), + size_bytes: 7, + }, + analysis_profile: AnalysisProfileManifest { + sha256: profile.clone(), + ghidra_version: "12.1.2".to_owned(), + adapter_protocol_version: 1, + target: BTreeMap::new(), + analyzer_options_sha256: digest('c'), + }, + created_by: CreatedBy { + ghidr_version: "0.1.0".to_owned(), + invocation_id: "0".repeat(32), + completed_at: "2026-01-01T00:00:00Z".to_owned(), + }, + project: AnalysisProject { + name: "analysis".to_owned(), + path: "project".to_owned(), + files: vec![InventoryFile { + path: "analysis.gpr".to_owned(), + size_bytes: 7, + sha256: hash_file(&project_file).expect("hash"), + }], + }, + }; + fs::write( + root.join("manifest.json"), + serde_json::to_vec(&manifest).expect("json"), + ) + .expect("manifest"); + } + + use std::path::Path; + + #[test] + fn complete_inventory_validates_and_mutation_is_corruption() { + let temp = TempDir::new().expect("temp"); + let sample = digest('a'); + let profile = digest('b'); + write_analysis(temp.path(), &sample, &profile); + validate_analysis(temp.path(), &sample, &profile).expect("valid"); + fs::write(temp.path().join("project/analysis.gpr"), b"changed").expect("change"); + assert!(validate_analysis(temp.path(), &sample, &profile).is_err()); + } + + #[cfg(unix)] + #[test] + fn inventory_rejects_symlink_attack() { + use std::os::unix::fs::symlink; + let temp = TempDir::new().expect("temp"); + let sample = digest('a'); + let profile = digest('b'); + write_analysis(temp.path(), &sample, &profile); + symlink("analysis.gpr", temp.path().join("project/escape")).expect("symlink"); + assert!(validate_analysis(temp.path(), &sample, &profile).is_err()); + } + + #[test] + fn recovery_allows_exactly_one_rebuild() { + assert_eq!(rebuild_decision(false), RebuildDecision::RebuildOnce); + assert_eq!(rebuild_decision(true), RebuildDecision::FailCorrupt); + } + + #[test] + fn promotion_is_atomic_and_corruption_moves_to_quarantine() { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(temp.path().join("store").into_os_string()), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let sample = digest('a'); + let profile = digest('b'); + let staged = store.root().join("staging/build/analysis"); + fs::create_dir_all(&staged).expect("staging"); + write_analysis(&staged, &sample, &profile); + let lock = AnalysisLock::acquire(&store, &sample, &profile).expect("lock"); + let active = store + .promote_analysis(&staged, &sample, &profile, &lock) + .expect("promote"); + assert!(!staged.exists()); + validate_analysis(&active, &sample, &profile).expect("active"); + fs::write(active.join("project/analysis.gpr"), b"corrupt").expect("corrupt"); + assert!(validate_analysis(&active, &sample, &profile).is_err()); + let outcome = store + .quarantine_analysis(&sample, &profile, &lock, false) + .expect("quarantine"); + assert!(!active.exists()); + assert!(outcome.path.exists()); + assert_eq!(outcome.decision, RebuildDecision::RebuildOnce); + } +} diff --git a/src/store/artifact.rs b/src/store/artifact.rs new file mode 100644 index 0000000..b46c596 --- /dev/null +++ b/src/store/artifact.rs @@ -0,0 +1,196 @@ +use std::{ + fs, + io::Write as _, + path::{Path, PathBuf}, +}; + +use sha2::{Digest as _, Sha256}; + +use crate::{ + domain::{Digest, TaggedPath}, + operation::ArtifactDescriptor, +}; + +use super::{ + AnalysisStore, StoreError, StoreErrorKind, + analysis::{digest_bytes, hash_file}, + filesystem::{ + create_private_directory, open_new_private_file, random_id, set_file_mode, sync_directory, + }, +}; + +/// Published immutable complete-success Artifact. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct Artifact { + path: PathBuf, + bytes: u64, + digest: Digest, +} + +impl Artifact { + /// Absolute published JSON path. + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + /// Exact bytes including the required trailing LF. + #[must_use] + pub const fn bytes(&self) -> u64 { + self.bytes + } + + /// SHA-256 of the exact published bytes. + #[must_use] + pub const fn digest(&self) -> &Digest { + &self.digest + } + + /// Builds the fixed public spill descriptor for this Artifact. + #[must_use] + pub fn descriptor(&self) -> ArtifactDescriptor { + ArtifactDescriptor { + path: TaggedPath::from_path(&self.path), + bytes: self.bytes, + sha256: self.digest.clone(), + media_type: "application/json".to_owned(), + contains: "complete_success_response".to_owned(), + } + } +} + +/// Atomically publishes an exact complete JSON success document by content digest. +pub fn publish_artifact( + store: &AnalysisStore, + sample: &Digest, + profile: &Digest, + exact_bytes: &[u8], +) -> Result { + validate_artifact_bytes(exact_bytes)?; + let digest = digest_bytes(&Sha256::digest(exact_bytes))?; + let directory = store + .path("artifacts") + .join(sample.as_str()) + .join(profile.as_str()); + create_private_directory(&directory)?; + let destination = directory.join(format!("{}.json", digest.as_str())); + if let Ok(metadata) = fs::symlink_metadata(&destination) { + if !metadata.is_file() + || metadata.file_type().is_symlink() + || hash_file(&destination)? != digest + { + return Err(StoreError::at( + StoreErrorKind::ImmutableConflict, + &destination, + "Artifact identity conflicts with existing storage", + )); + } + return artifact(destination, exact_bytes.len(), digest); + } + + let temporary = directory.join(format!("artifact-{}.tmp", random_id()?)); + let result = (|| { + let mut file = open_new_private_file(&temporary)?; + file.write_all(exact_bytes) + .map_err(|error| StoreError::io(&temporary, error))?; + file.sync_all() + .map_err(|error| StoreError::io(&temporary, error))?; + set_file_mode(&temporary, 0o400)?; + match fs::hard_link(&temporary, &destination) { + Ok(()) => {} + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => { + if hash_file(&destination)? != digest { + return Err(StoreError::at( + StoreErrorKind::ImmutableConflict, + &destination, + "concurrent Artifact conflict", + )); + } + } + Err(error) => return Err(StoreError::io(&destination, error)), + } + fs::remove_file(&temporary).map_err(|error| StoreError::io(&temporary, error))?; + sync_directory(&directory) + })(); + if result.is_err() { + let _ = fs::remove_file(&temporary); + } + result?; + artifact(destination, exact_bytes.len(), digest) +} + +fn validate_artifact_bytes(bytes: &[u8]) -> Result<(), StoreError> { + if bytes.is_empty() || !bytes.ends_with(b"\n") || bytes.ends_with(b"\n\n") { + return Err(StoreError::new( + StoreErrorKind::Io, + "Artifact must have exactly one trailing LF", + )); + } + let json = &bytes[..bytes.len() - 1]; + if json.starts_with(&[0xef, 0xbb, 0xbf]) + || std::str::from_utf8(json).is_err() + || serde_json::from_slice::(json).is_err() + { + return Err(StoreError::new( + StoreErrorKind::Io, + "Artifact must be BOM-free valid UTF-8 JSON", + )); + } + Ok(()) +} + +fn artifact(path: PathBuf, length: usize, digest: Digest) -> Result { + Ok(Artifact { + path, + bytes: u64::try_from(length) + .map_err(|_| StoreError::new(StoreErrorKind::Io, "Artifact length overflow"))?, + digest, + }) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use super::publish_artifact; + use crate::{ + domain::Digest, + store::{AnalysisStore, StoreEnvironment, resolve_store}, + }; + use std::{ffi::OsString, fs, str::FromStr as _}; + use tempfile::TempDir; + + #[test] + fn identical_exact_bytes_publish_one_immutable_identity() { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let sample = Digest::from_str(&"a".repeat(64)).expect("digest"); + let profile = Digest::from_str(&"b".repeat(64)).expect("digest"); + let first = publish_artifact(&store, &sample, &profile, b"{\"kind\":\"inspection\"}\n") + .expect("publish"); + let second = publish_artifact(&store, &sample, &profile, b"{\"kind\":\"inspection\"}\n") + .expect("reuse"); + assert_eq!(first, second); + assert_eq!( + fs::read(first.path()).expect("read"), + b"{\"kind\":\"inspection\"}\n" + ); + } + + #[test] + fn refuses_partial_or_non_json_artifacts() { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let digest = Digest::from_str(&"a".repeat(64)).expect("digest"); + assert!(publish_artifact(&store, &digest, &digest, b"{}").is_err()); + } +} diff --git a/src/store/cleanup.rs b/src/store/cleanup.rs new file mode 100644 index 0000000..e409f1b --- /dev/null +++ b/src/store/cleanup.rs @@ -0,0 +1,738 @@ +use std::{ + collections::BTreeSet, + fs, + path::{Path, PathBuf}, + str::FromStr as _, +}; + +use serde::{Deserialize, Serialize}; + +use crate::{ + domain::Digest, + operation::{CleanupSnapshot, StorageUsage}, +}; + +use super::{ + AnalysisLock, AnalysisStore, StoreError, StoreErrorKind, + filesystem::{ + atomic_json_write, create_private_directory, random_id, scan_usage, sync_directory, + }, +}; + +/// Store-wide or one-Sample cleanup selection. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum CleanupScope { + /// All data tied to this Sample digest. + Sample(Digest), + /// Every removable immutable object and Diagnostic Log. + All, +} + +/// Deterministic non-mutating cleanup preflight, including resumable staged data. +#[derive(Debug)] +pub struct CleanupPlan { + scope: CleanupScope, + entries: Vec, + snapshot: CleanupSnapshot, +} + +impl CleanupPlan { + /// Exact counts, usage, and sorted Profile identities selected by preflight. + #[must_use] + pub const fn snapshot(&self) -> &CleanupSnapshot { + &self.snapshot + } + + /// Executes locking, atomic staging, deletion, and transaction recovery. + pub fn execute(self, store: &AnalysisStore) -> Result { + execute_plan(store, self) + } +} + +/// Result of an executed cleanup transaction. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CleanupTransaction { + /// Exact preflight snapshot. + pub matched: CleanupSnapshot, + /// Exact successfully removed subset. + pub removed: CleanupSnapshot, + /// Retained transaction path only when deletion is incomplete. + pub transaction_path: Option, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(rename_all = "snake_case")] +enum Category { + Analysis, + Quarantine, + Artifact, + Diagnostic, +} + +#[derive(Clone, Debug)] +enum EntryState { + Active, + Staged, +} + +#[derive(Clone, Debug)] +struct PlanEntry { + path: PathBuf, + original: PathBuf, + state: EntryState, + category: Category, + sample: Option, + profile: Option, + usage: StorageUsage, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +struct TransactionManifest { + transaction_version: u32, + scope: TransactionScope, + entries: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)] +#[serde(tag = "kind", content = "sha256", rename_all = "snake_case")] +enum TransactionScope { + Sample(Digest), + All, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +struct TransactionEntry { + staged_name: String, + original_path: String, + category: Category, + sample: Option, + profile: Option, + usage: StorageUsage, +} + +/// Traverses active and previously staged objects without changing store data. +pub fn plan_cleanup(store: &AnalysisStore, scope: CleanupScope) -> Result { + let mut entries = active_entries(store, &scope)?; + entries.extend(staged_entries(store, &scope)?); + entries.sort_by(|left, right| path_bytes(&left.path).cmp(path_bytes(&right.path))); + let snapshot = summarize(&entries)?; + Ok(CleanupPlan { + scope, + entries, + snapshot, + }) +} + +fn active_entries( + store: &AnalysisStore, + scope: &CleanupScope, +) -> Result, StoreError> { + let mut entries = Vec::new(); + collect_profile_directories( + store, + scope, + "analyses", + Category::Analysis, + false, + &mut entries, + )?; + collect_profile_directories( + store, + scope, + "quarantine", + Category::Quarantine, + true, + &mut entries, + )?; + collect_artifacts(store, scope, &mut entries)?; + collect_diagnostics(store, scope, &mut entries)?; + Ok(entries) +} + +fn collect_profile_directories( + store: &AnalysisStore, + scope: &CleanupScope, + top: &str, + category: Category, + has_invocations: bool, + entries: &mut Vec, +) -> Result<(), StoreError> { + for sample_entry in selected_sample_directories(&store.path(top), scope)? { + let sample = parse_digest_name(&sample_entry)?; + for profile_entry in directories(&sample_entry)? { + let profile = parse_digest_name(&profile_entry)?; + if has_invocations { + for invocation in directories(&profile_entry)? { + push_active( + entries, + invocation, + category, + Some(sample.clone()), + Some(profile.clone()), + )?; + } + } else { + push_active( + entries, + profile_entry, + category, + Some(sample.clone()), + Some(profile), + )?; + } + } + } + Ok(()) +} + +fn collect_artifacts( + store: &AnalysisStore, + scope: &CleanupScope, + entries: &mut Vec, +) -> Result<(), StoreError> { + for sample_entry in selected_sample_directories(&store.path("artifacts"), scope)? { + let sample = parse_digest_name(&sample_entry)?; + for profile_entry in directories(&sample_entry)? { + let profile = parse_digest_name(&profile_entry)?; + for artifact in regular_files(&profile_entry)? { + push_active( + entries, + artifact, + Category::Artifact, + Some(sample.clone()), + Some(profile.clone()), + )?; + } + } + } + Ok(()) +} + +fn collect_diagnostics( + store: &AnalysisStore, + scope: &CleanupScope, + entries: &mut Vec, +) -> Result<(), StoreError> { + let root = store.path("diagnostics"); + for directory in directories(&root)? { + let name = file_name_utf8(&directory)?; + let sample = if name == "global" { + if !matches!(scope, CleanupScope::All) { + continue; + } + None + } else { + let digest = parse_digest_name(&directory)?; + if !scope_matches(scope, &digest) { + continue; + } + Some(digest) + }; + for bundle in directories(&directory)? { + push_active(entries, bundle, Category::Diagnostic, sample.clone(), None)?; + } + } + Ok(()) +} + +fn push_active( + entries: &mut Vec, + path: PathBuf, + category: Category, + sample: Option, + profile: Option, +) -> Result<(), StoreError> { + let usage = scan_usage(&path)?; + entries.push(PlanEntry { + original: path.clone(), + path, + state: EntryState::Active, + category, + sample, + profile, + usage, + }); + Ok(()) +} + +fn staged_entries( + store: &AnalysisStore, + scope: &CleanupScope, +) -> Result, StoreError> { + let mut result = Vec::new(); + for transaction in directories(&store.path("cleanup"))? { + let manifest_path = transaction.join("manifest.json"); + let bytes = + fs::read(&manifest_path).map_err(|error| StoreError::io(&manifest_path, error))?; + let manifest: TransactionManifest = serde_json::from_slice(&bytes).map_err(|error| { + StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + format!("invalid cleanup manifest: {error}"), + ) + })?; + if manifest.transaction_version != 1 || !transaction_scope_matches(&manifest.scope, scope) { + continue; + } + for entry in manifest.entries { + let path = transaction.join("payload").join(&entry.staged_name); + match fs::symlink_metadata(&path) { + Ok(_) => result.push(PlanEntry { + path, + original: PathBuf::from(entry.original_path), + state: EntryState::Staged, + category: entry.category, + sample: entry.sample, + profile: entry.profile, + usage: entry.usage, + }), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(StoreError::io(&path, error)), + } + } + } + Ok(result) +} + +fn execute_plan( + store: &AnalysisStore, + plan: CleanupPlan, +) -> Result { + let mut lock_keys: BTreeSet<(Digest, Digest)> = BTreeSet::new(); + for entry in &plan.entries { + if let (Some(sample), Some(profile)) = (&entry.sample, &entry.profile) { + lock_keys.insert((sample.clone(), profile.clone())); + } + } + let mut locks = Vec::new(); + for (sample, profile) in lock_keys { + locks.push(AnalysisLock::acquire(store, &sample, &profile)?); + } + + let staged_existing: Vec<_> = plan + .entries + .iter() + .filter(|entry| matches!(entry.state, EntryState::Staged)) + .cloned() + .collect(); + let active: Vec<_> = plan + .entries + .iter() + .filter(|entry| matches!(entry.state, EntryState::Active)) + .cloned() + .collect(); + let mut removed_entries = Vec::new(); + for entry in staged_existing { + remove_entry(&entry.path).map_err(|error| cleanup_incomplete(&entry.path, error))?; + removed_entries.push(entry); + } + remove_empty_transactions(store)?; + + let transaction_path = if active.is_empty() { + None + } else { + Some(stage_active(store, &plan.scope, &active)?) + }; + if let Some(transaction) = &transaction_path { + let staged = staged_entries_in_transaction(transaction)?; + for entry in staged { + if let Err(error) = remove_entry(&entry.path) { + let removed = summarize(&removed_entries)?; + return Err(StoreError::at( + StoreErrorKind::CleanupIncomplete, + transaction, + format!( + "cleanup deletion incomplete after removing {} objects: {error}", + count_objects(&removed) + ), + )); + } + removed_entries.push(entry); + } + fs::remove_dir_all(transaction).map_err(|error| cleanup_incomplete(transaction, error))?; + sync_directory(&store.path("cleanup"))?; + } + drop(locks); + Ok(CleanupTransaction { + matched: plan.snapshot, + removed: summarize(&removed_entries)?, + transaction_path: None, + }) +} + +fn stage_active( + store: &AnalysisStore, + scope: &CleanupScope, + entries: &[PlanEntry], +) -> Result { + let transaction = store.path("cleanup").join(random_id()?); + let payload = transaction.join("payload"); + create_private_directory(&payload)?; + let manifest_entries: Vec<_> = entries + .iter() + .enumerate() + .map(|(index, entry)| TransactionEntry { + staged_name: format!("{index:016x}"), + original_path: entry.original.to_string_lossy().into_owned(), + category: entry.category, + sample: entry.sample.clone(), + profile: entry.profile.clone(), + usage: entry.usage, + }) + .collect(); + let manifest = TransactionManifest { + transaction_version: 1, + scope: match scope { + CleanupScope::Sample(digest) => TransactionScope::Sample(digest.clone()), + CleanupScope::All => TransactionScope::All, + }, + entries: manifest_entries, + }; + atomic_json_write(&transaction, &transaction.join("manifest.json"), &manifest)?; + let mut completed = Vec::new(); + for (index, entry) in entries.iter().enumerate() { + let staged = payload.join(format!("{index:016x}")); + if let Err(error) = fs::rename(&entry.path, &staged) { + for (original, staged) in completed.into_iter().rev() { + let _ = fs::rename(staged, original); + } + let _ = fs::remove_dir_all(&transaction); + return Err(StoreError::io(&entry.path, error)); + } + completed.push((entry.original.clone(), staged)); + } + sync_directory(&payload)?; + Ok(transaction) +} + +fn staged_entries_in_transaction(transaction: &Path) -> Result, StoreError> { + let manifest_path = transaction.join("manifest.json"); + let manifest: TransactionManifest = serde_json::from_slice( + &fs::read(&manifest_path).map_err(|error| StoreError::io(&manifest_path, error))?, + ) + .map_err(|error| { + StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + format!("invalid cleanup manifest: {error}"), + ) + })?; + Ok(manifest + .entries + .into_iter() + .map(|entry| PlanEntry { + path: transaction.join("payload").join(entry.staged_name), + original: PathBuf::from(entry.original_path), + state: EntryState::Staged, + category: entry.category, + sample: entry.sample, + profile: entry.profile, + usage: entry.usage, + }) + .collect()) +} + +fn remove_entry(path: &Path) -> std::io::Result<()> { + let metadata = fs::symlink_metadata(path)?; + if metadata.is_dir() { + fs::remove_dir_all(path) + } else { + fs::remove_file(path) + } +} + +fn remove_empty_transactions(store: &AnalysisStore) -> Result<(), StoreError> { + for transaction in directories(&store.path("cleanup"))? { + let payload = transaction.join("payload"); + if directories_and_files(&payload)?.is_empty() { + fs::remove_dir_all(&transaction) + .map_err(|error| StoreError::io(&transaction, error))?; + } + } + Ok(()) +} + +fn summarize(entries: &[PlanEntry]) -> Result { + let mut result = CleanupSnapshot { + analyses: 0, + quarantined_analyses: 0, + artifacts: 0, + diagnostic_logs: 0, + usage: StorageUsage { + logical_bytes: 0, + allocated_bytes: 0, + }, + analysis_profile_sha256: Vec::new(), + }; + let mut profiles = BTreeSet::new(); + for entry in entries { + match entry.category { + Category::Analysis => result.analyses += 1, + Category::Quarantine => result.quarantined_analyses += 1, + Category::Artifact => result.artifacts += 1, + Category::Diagnostic => result.diagnostic_logs += 1, + } + result.usage.logical_bytes = result + .usage + .logical_bytes + .checked_add(entry.usage.logical_bytes) + .ok_or_else(|| { + StoreError::new( + StoreErrorKind::CorruptStore, + "cleanup logical usage overflow", + ) + })?; + result.usage.allocated_bytes = result + .usage + .allocated_bytes + .checked_add(entry.usage.allocated_bytes) + .ok_or_else(|| { + StoreError::new( + StoreErrorKind::CorruptStore, + "cleanup allocated usage overflow", + ) + })?; + if let Some(profile) = &entry.profile { + profiles.insert(profile.clone()); + } + } + result.analysis_profile_sha256 = profiles.into_iter().collect(); + Ok(result) +} + +fn selected_sample_directories( + root: &Path, + scope: &CleanupScope, +) -> Result, StoreError> { + let directories = directories(root)?; + directories + .into_iter() + .filter_map(|path| match parse_digest_name(&path) { + Ok(digest) if scope_matches(scope, &digest) => Some(Ok(path)), + Ok(_) => None, + Err(error) => Some(Err(error)), + }) + .collect() +} + +fn directories(path: &Path) -> Result, StoreError> { + filter_entries(path, true) +} + +fn regular_files(path: &Path) -> Result, StoreError> { + filter_entries(path, false) +} + +fn directories_and_files(path: &Path) -> Result, StoreError> { + let mut result = Vec::new(); + for entry in read_entries(path)? { + let metadata = fs::symlink_metadata(entry.path()) + .map_err(|error| StoreError::io(entry.path(), error))?; + if metadata.file_type().is_symlink() || !(metadata.is_dir() || metadata.is_file()) { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + entry.path(), + "unsafe cleanup entry", + )); + } + result.push(entry.path()); + } + Ok(result) +} + +fn filter_entries(path: &Path, want_directory: bool) -> Result, StoreError> { + let mut result = Vec::new(); + for entry in read_entries(path)? { + let metadata = fs::symlink_metadata(entry.path()) + .map_err(|error| StoreError::io(entry.path(), error))?; + if metadata.file_type().is_symlink() + || metadata.is_dir() != want_directory + || metadata.is_file() == want_directory + { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + entry.path(), + "unexpected or unsafe store entry", + )); + } + result.push(entry.path()); + } + Ok(result) +} + +fn read_entries(path: &Path) -> Result, StoreError> { + let mut entries = fs::read_dir(path) + .map_err(|error| StoreError::io(path, error))? + .collect::, _>>() + .map_err(|error| StoreError::io(path, error))?; + entries.sort_by_key(fs::DirEntry::file_name); + Ok(entries) +} + +fn parse_digest_name(path: &Path) -> Result { + Digest::from_str(file_name_utf8(path)?).map_err(|_| { + StoreError::at( + StoreErrorKind::CorruptStore, + path, + "non-digest store component", + ) + }) +} + +fn file_name_utf8(path: &Path) -> Result<&str, StoreError> { + path.file_name() + .and_then(std::ffi::OsStr::to_str) + .ok_or_else(|| { + StoreError::at( + StoreErrorKind::CorruptStore, + path, + "non-UTF-8 generated component", + ) + }) +} + +fn scope_matches(scope: &CleanupScope, digest: &Digest) -> bool { + matches!(scope, CleanupScope::All) + || matches!(scope, CleanupScope::Sample(selected) if selected == digest) +} + +fn transaction_scope_matches(transaction: &TransactionScope, scope: &CleanupScope) -> bool { + matches!( + (transaction, scope), + (TransactionScope::All, CleanupScope::All) + ) || matches!((transaction, scope), (TransactionScope::Sample(left), CleanupScope::Sample(right)) if left == right) + || matches!(scope, CleanupScope::All) +} + +#[cfg(unix)] +fn path_bytes(path: &Path) -> &[u8] { + use std::os::unix::ffi::OsStrExt as _; + path.as_os_str().as_bytes() +} + +fn cleanup_incomplete(path: &Path, error: std::io::Error) -> StoreError { + StoreError::at( + StoreErrorKind::CleanupIncomplete, + path, + format!("cleanup deletion incomplete: {error}"), + ) +} + +fn count_objects(snapshot: &CleanupSnapshot) -> u64 { + snapshot.analyses + + snapshot.quarantined_analyses + + snapshot.artifacts + + snapshot.diagnostic_logs +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use super::{CleanupScope, plan_cleanup}; + use crate::{ + domain::Digest, + store::{AnalysisLock, AnalysisStore, StoreEnvironment, resolve_store}, + }; + use std::{ffi::OsString, fs, str::FromStr as _}; + use tempfile::TempDir; + + fn setup() -> (TempDir, AnalysisStore, Digest, Digest) { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let sample = Digest::from_str(&"a".repeat(64)).expect("digest"); + let profile = Digest::from_str(&"b".repeat(64)).expect("digest"); + (temp, store, sample, profile) + } + + #[test] + fn dry_run_plan_is_exact_and_non_mutating() { + let (_temp, store, sample, profile) = setup(); + let path = store + .root() + .join("artifacts") + .join(sample.as_str()) + .join(profile.as_str()); + fs::create_dir_all(&path).expect("dirs"); + fs::write(path.join(format!("{}.json", "c".repeat(64))), b"{}\n").expect("artifact"); + let plan = plan_cleanup(&store, CleanupScope::Sample(sample)).expect("plan"); + assert_eq!(plan.snapshot().artifacts, 1); + assert!(path.exists()); + } + + #[test] + fn execution_stages_then_removes_selected_objects_only() { + let (_temp, store, sample, profile) = setup(); + let selected = store + .root() + .join("artifacts") + .join(sample.as_str()) + .join(profile.as_str()); + fs::create_dir_all(&selected).expect("dirs"); + fs::write(selected.join(format!("{}.json", "c".repeat(64))), b"{}\n").expect("artifact"); + let other = Digest::from_str(&"d".repeat(64)).expect("digest"); + let other_path = store + .root() + .join("artifacts") + .join(other.as_str()) + .join(profile.as_str()); + fs::create_dir_all(&other_path).expect("dirs"); + fs::write(other_path.join(format!("{}.json", "e".repeat(64))), b"{}\n").expect("artifact"); + let result = plan_cleanup(&store, CleanupScope::Sample(sample)) + .expect("plan") + .execute(&store) + .expect("execute"); + assert_eq!(result.removed.artifacts, 1); + assert!(!selected.join(format!("{}.json", "c".repeat(64))).exists()); + assert!(other_path.exists()); + } + + #[test] + fn staged_transaction_is_discovered_for_resumable_recovery() { + let (_temp, store, sample, profile) = setup(); + let transaction = store.root().join("cleanup").join("0".repeat(32)); + fs::create_dir_all(transaction.join("payload")).expect("payload"); + fs::write(transaction.join("payload/0000000000000000"), b"{}\n").expect("staged"); + let manifest = serde_json::json!({ + "transaction_version": 1, + "scope": {"kind":"sample", "sha256": sample}, + "entries": [{"staged_name":"0000000000000000", "original_path":"/gone", "category":"artifact", "sample": sample, "profile": profile, "usage":{"logical_bytes":3,"allocated_bytes":4096}}] + }); + fs::write( + transaction.join("manifest.json"), + serde_json::to_vec(&manifest).expect("json"), + ) + .expect("manifest"); + let plan = plan_cleanup( + &store, + CleanupScope::Sample(Digest::from_str(&"a".repeat(64)).expect("digest")), + ) + .expect("plan"); + assert_eq!(plan.snapshot().artifacts, 1); + plan.execute(&store).expect("resume"); + assert!(!transaction.exists()); + } + + #[test] + fn busy_analysis_aborts_before_store_mutation() { + let (_temp, store, sample, profile) = setup(); + let analysis = store + .root() + .join("analyses") + .join(sample.as_str()) + .join(profile.as_str()); + fs::create_dir_all(&analysis).expect("analysis"); + fs::write(analysis.join("data"), b"active").expect("data"); + let held = AnalysisLock::acquire(&store, &sample, &profile).expect("lock"); + let plan = plan_cleanup(&store, CleanupScope::Sample(sample)).expect("plan"); + assert!(plan.execute(&store).is_err()); + assert!(analysis.exists()); + drop(held); + } +} diff --git a/src/store/cleanup_command.rs b/src/store/cleanup_command.rs new file mode 100644 index 0000000..24a71c7 --- /dev/null +++ b/src/store/cleanup_command.rs @@ -0,0 +1,215 @@ +use serde_json::Value; + +use crate::{ + AppError, ErrorCode, + cli::CleanArgs, + domain::TaggedPath, + operation::{CleanupData, CleanupMode, CleanupTarget}, +}; + +use super::{ + AnalysisStore, CleanupScope, SnapshotOptions, StoreError, StoreErrorKind, plan_cleanup, + snapshot_sample, +}; + +/// Executes the store-facing portion of the public `clean` command synchronously. +pub fn execute_cleanup( + store: &AnalysisStore, + arguments: &CleanArgs, + max_sample_bytes: u64, +) -> Result { + let (scope, target) = if let Some(sample) = &arguments.sample { + let snapshot = snapshot_sample( + store, + sample, + SnapshotOptions { + max_sample_bytes, + reserve_bytes: 0, + }, + ) + .map_err(app_error)?; + ( + CleanupScope::Sample(snapshot.digest().clone()), + CleanupTarget::Sample { + source_path: TaggedPath::from_path(sample), + sha256: snapshot.digest().clone(), + }, + ) + } else if let Some(digest) = &arguments.digest { + ( + CleanupScope::Sample(digest.clone()), + CleanupTarget::Digest { + sha256: digest.clone(), + }, + ) + } else if arguments.all { + (CleanupScope::All, CleanupTarget::All) + } else { + return Err(AppError::invalid_arguments( + "exactly one cleanup target is required", + )); + }; + + let plan = plan_cleanup(store, scope.clone()).map_err(app_error)?; + if arguments.all && !arguments.dry_run && !arguments.yes { + return Err(AppError::new( + ErrorCode::ConfirmationRequired, + "store-wide cleanup requires explicit confirmation", + false, + ) + .with_detail("analyses", Value::from(plan.snapshot().analyses)) + .with_detail("artifacts", Value::from(plan.snapshot().artifacts)) + .with_detail( + "logical_bytes", + Value::from(plan.snapshot().usage.logical_bytes), + ) + .with_detail( + "allocated_bytes", + Value::from(plan.snapshot().usage.allocated_bytes), + )); + } + if arguments.dry_run { + return Ok(CleanupData { + mode: CleanupMode::DryRun, + target, + matched: plan.snapshot().clone(), + removed: None, + }); + } + let matched = plan.snapshot().clone(); + let transaction = match plan.execute(store) { + Ok(transaction) => transaction, + Err(error) if error.kind() == StoreErrorKind::CleanupIncomplete => { + let remaining = plan_cleanup(store, scope) + .map(|plan| plan.snapshot().clone()) + .unwrap_or_else(|_| matched.clone()); + let removed = subtract_snapshot(&matched, &remaining); + let mut app = app_error(error); + app = app.with_detail( + "removed", + serde_json::to_value(removed).unwrap_or(Value::Null), + ); + app = app.with_detail( + "remaining", + serde_json::to_value(remaining).unwrap_or(Value::Null), + ); + return Err(app); + } + Err(error) => return Err(app_error(error)), + }; + Ok(CleanupData { + mode: CleanupMode::Executed, + target, + matched: transaction.matched, + removed: Some(transaction.removed), + }) +} + +fn subtract_snapshot( + matched: &crate::operation::CleanupSnapshot, + remaining: &crate::operation::CleanupSnapshot, +) -> crate::operation::CleanupSnapshot { + use crate::operation::{CleanupSnapshot, StorageUsage}; + + CleanupSnapshot { + analyses: matched.analyses.saturating_sub(remaining.analyses), + quarantined_analyses: matched + .quarantined_analyses + .saturating_sub(remaining.quarantined_analyses), + artifacts: matched.artifacts.saturating_sub(remaining.artifacts), + diagnostic_logs: matched + .diagnostic_logs + .saturating_sub(remaining.diagnostic_logs), + usage: StorageUsage { + logical_bytes: matched + .usage + .logical_bytes + .saturating_sub(remaining.usage.logical_bytes), + allocated_bytes: matched + .usage + .allocated_bytes + .saturating_sub(remaining.usage.allocated_bytes), + }, + analysis_profile_sha256: matched + .analysis_profile_sha256 + .iter() + .filter(|profile| !remaining.analysis_profile_sha256.contains(profile)) + .cloned() + .collect(), + } +} + +fn app_error(error: StoreError) -> AppError { + let (code, retryable) = match error.kind() { + StoreErrorKind::InvalidStorePath | StoreErrorKind::UnsafeStorePath => { + (ErrorCode::InvalidStorePath, false) + } + StoreErrorKind::CorruptStore | StoreErrorKind::ImmutableConflict => { + (ErrorCode::CorruptAnalysis, false) + } + StoreErrorKind::SampleNotFound => (ErrorCode::SampleNotFound, false), + StoreErrorKind::SampleUnreadable => (ErrorCode::SampleUnreadable, false), + StoreErrorKind::InvalidSampleType => (ErrorCode::InvalidSampleType, false), + StoreErrorKind::SampleTooLarge => (ErrorCode::SampleTooLarge, false), + StoreErrorKind::SampleChanged => (ErrorCode::SampleChanged, true), + StoreErrorKind::InsufficientStoreSpace => (ErrorCode::InsufficientStoreSpace, true), + StoreErrorKind::AnalysisBusy => (ErrorCode::AnalysisBusy, true), + StoreErrorKind::CleanupIncomplete => (ErrorCode::CleanupIncomplete, true), + StoreErrorKind::Io => (ErrorCode::Internal, false), + }; + let message = error.to_string(); + let mut app_error = AppError::new(code, message, retryable); + if let Some(path) = error.path() { + app_error = app_error.with_detail( + "path", + serde_json::to_value(TaggedPath::from_path(path)).unwrap_or(Value::Null), + ); + } + app_error +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{ffi::OsString, fs, str::FromStr as _}; + use tempfile::TempDir; + + use super::execute_cleanup; + use crate::{ + cli::CleanArgs, + domain::Digest, + operation::CleanupMode, + store::{AnalysisStore, StoreEnvironment, resolve_store}, + }; + + #[test] + fn dry_run_command_reports_without_deleting() { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let sample = Digest::from_str(&"a".repeat(64)).expect("digest"); + let profile = Digest::from_str(&"b".repeat(64)).expect("digest"); + let directory = store + .root() + .join("artifacts") + .join(sample.as_str()) + .join(profile.as_str()); + fs::create_dir_all(&directory).expect("dirs"); + fs::write(directory.join(format!("{}.json", "c".repeat(64))), b"{}\n").expect("artifact"); + let args = CleanArgs { + sample: None, + digest: Some(sample), + all: false, + dry_run: true, + yes: false, + }; + let data = execute_cleanup(&store, &args, 1024).expect("clean"); + assert_eq!(data.mode, CleanupMode::DryRun); + assert_eq!(data.matched.artifacts, 1); + assert!(directory.exists()); + } +} diff --git a/src/store/error.rs b/src/store/error.rs new file mode 100644 index 0000000..a5ef133 --- /dev/null +++ b/src/store/error.rs @@ -0,0 +1,96 @@ +use std::{io, path::PathBuf}; + +use thiserror::Error; + +/// Stable store-layer failure classification used by command integration. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum StoreErrorKind { + /// A configured store path violates the store path contract. + InvalidStorePath, + /// A path inside tool-owned storage is a symlink or otherwise unsafe. + UnsafeStorePath, + /// The store layout or an immutable object is corrupt. + CorruptStore, + /// The Sample is absent. + SampleNotFound, + /// The Sample could not be read. + SampleUnreadable, + /// The opened Sample is not a regular file. + InvalidSampleType, + /// The Sample exceeds the configured streaming bound. + SampleTooLarge, + /// The source changed while the private snapshot was copied. + SampleChanged, + /// The store filesystem does not have the required reserve. + InsufficientStoreSpace, + /// Another process owns the requested Analysis lock. + AnalysisBusy, + /// An immutable destination already exists with different content. + ImmutableConflict, + /// A cleanup transaction was staged but could not be fully deleted. + CleanupIncomplete, + /// General filesystem or serialization failure. + Io, +} + +/// Detailed synchronous store failure. +#[derive(Debug, Error)] +#[error("{message}")] +pub struct StoreError { + kind: StoreErrorKind, + message: String, + path: Option, + source: Option, +} + +impl StoreError { + pub(crate) fn new(kind: StoreErrorKind, message: impl Into) -> Self { + Self { + kind, + message: message.into(), + path: None, + source: None, + } + } + + pub(crate) fn at( + kind: StoreErrorKind, + path: impl Into, + message: impl Into, + ) -> Self { + Self { + kind, + message: message.into(), + path: Some(path.into()), + source: None, + } + } + + pub(crate) fn io(path: impl Into, source: io::Error) -> Self { + let path = path.into(); + Self { + kind: StoreErrorKind::Io, + message: format!("filesystem operation failed for {}", path.display()), + path: Some(path), + source: Some(source), + } + } + + /// Returns the stable failure class. + #[must_use] + pub const fn kind(&self) -> StoreErrorKind { + self.kind + } + + /// Returns the exact implicated path when available. + #[must_use] + pub fn path(&self) -> Option<&std::path::Path> { + self.path.as_deref() + } + + /// Returns the underlying I/O error when one was retained. + #[must_use] + pub const fn io_source(&self) -> Option<&io::Error> { + self.source.as_ref() + } +} diff --git a/src/store/filesystem.rs b/src/store/filesystem.rs new file mode 100644 index 0000000..b94172e --- /dev/null +++ b/src/store/filesystem.rs @@ -0,0 +1,332 @@ +use std::{ + fs::{self, File, OpenOptions}, + io::Write as _, + os::unix::fs::{MetadataExt as _, OpenOptionsExt as _, PermissionsExt as _}, + path::{Path, PathBuf}, +}; + +use serde::{Deserialize, Serialize}; + +use crate::operation::{StorageUsage, StoreSource}; + +use super::{DIRECTORY_MODE, FILE_MODE, ResolvedStore, StoreError, StoreErrorKind}; + +/// Current on-disk layout version recorded in `store.json`. +pub const STORE_LAYOUT_VERSION: u32 = 1; + +const LAYOUT_DIRECTORIES: &[&str] = &[ + "analyses", + "artifacts", + "diagnostics", + "quarantine", + "locks", + "staging", + "cleanup", +]; + +#[derive(Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +struct StoreManifest { + store_layout_version: u32, +} + +/// Initialized secure Analysis Store. +#[derive(Clone, Debug)] +pub struct AnalysisStore { + root: PathBuf, + source: StoreSource, +} + +impl AnalysisStore { + /// Creates or validates the version 1 layout with private permissions. + pub fn initialize(resolved: ResolvedStore) -> Result { + create_private_directory(resolved.path())?; + for component in LAYOUT_DIRECTORIES { + create_private_directory(&resolved.path().join(component))?; + } + let manifest_path = resolved.path().join("store.json"); + if manifest_path.exists() { + reject_symlink(&manifest_path)?; + let bytes = + fs::read(&manifest_path).map_err(|error| StoreError::io(&manifest_path, error))?; + let manifest: StoreManifest = serde_json::from_slice(&bytes).map_err(|error| { + StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + format!("store.json is invalid: {error}"), + ) + })?; + if manifest.store_layout_version != STORE_LAYOUT_VERSION { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + &manifest_path, + "unsupported Analysis Store layout version", + )); + } + set_file_mode(&manifest_path, FILE_MODE)?; + } else { + let manifest = StoreManifest { + store_layout_version: STORE_LAYOUT_VERSION, + }; + atomic_json_write(resolved.path(), &manifest_path, &manifest)?; + } + Ok(Self { + root: resolved.path().to_path_buf(), + source: resolved.source(), + }) + } + + /// Returns the absolute store root. + #[must_use] + pub fn root(&self) -> &Path { + &self.root + } + + /// Returns the resolution precedence source. + #[must_use] + pub const fn source(&self) -> StoreSource { + self.source + } + + /// Scans with `lstat`, rejecting symlinks and special files. + pub fn usage(&self) -> Result { + scan_usage(&self.root) + } + + pub(crate) fn path(&self, component: &str) -> PathBuf { + self.root.join(component) + } +} + +/// Logical file bytes and allocated blocks for a traversed tree. +pub type StoreUsage = StorageUsage; + +pub(crate) fn create_private_directory(path: &Path) -> Result<(), StoreError> { + reject_symlink_ancestors(path)?; + match fs::symlink_metadata(path) { + Ok(metadata) => { + if !metadata.file_type().is_dir() { + return Err(StoreError::at( + StoreErrorKind::UnsafeStorePath, + path, + "tool-owned path is not a directory", + )); + } + } + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + fs::create_dir_all(path).map_err(|error| StoreError::io(path, error))?; + reject_symlink(path)?; + } + Err(error) => return Err(StoreError::io(path, error)), + } + fs::set_permissions(path, fs::Permissions::from_mode(DIRECTORY_MODE)) + .map_err(|error| StoreError::io(path, error)) +} + +fn reject_symlink_ancestors(path: &Path) -> Result<(), StoreError> { + for ancestor in path.ancestors().collect::>().into_iter().rev() { + match fs::symlink_metadata(ancestor) { + Ok(metadata) if metadata.file_type().is_symlink() => { + return Err(StoreError::at( + StoreErrorKind::UnsafeStorePath, + ancestor, + "symlink found in tool-owned path traversal", + )); + } + Ok(metadata) if ancestor != path && !metadata.is_dir() => { + return Err(StoreError::at( + StoreErrorKind::UnsafeStorePath, + ancestor, + "non-directory found in tool-owned path traversal", + )); + } + Ok(_) => {} + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(StoreError::io(ancestor, error)), + } + } + Ok(()) +} + +pub(crate) fn reject_symlink(path: &Path) -> Result<(), StoreError> { + let metadata = fs::symlink_metadata(path).map_err(|error| StoreError::io(path, error))?; + if metadata.file_type().is_symlink() { + return Err(StoreError::at( + StoreErrorKind::UnsafeStorePath, + path, + "symlink found inside tool-owned storage", + )); + } + Ok(()) +} + +pub(crate) fn set_file_mode(path: &Path, mode: u32) -> Result<(), StoreError> { + fs::set_permissions(path, fs::Permissions::from_mode(mode)) + .map_err(|error| StoreError::io(path, error)) +} + +pub(crate) fn open_new_private_file(path: &Path) -> Result { + OpenOptions::new() + .write(true) + .create_new(true) + .mode(FILE_MODE) + .open(path) + .map_err(|error| StoreError::io(path, error)) +} + +pub(crate) fn atomic_json_write( + temporary_parent: &Path, + destination: &Path, + value: &T, +) -> Result<(), StoreError> { + let temporary = temporary_parent.join(format!("manifest-{}.tmp", random_id()?)); + let mut bytes = serde_json::to_vec(value).map_err(|error| { + StoreError::new( + StoreErrorKind::Io, + format!("JSON serialization failed: {error}"), + ) + })?; + bytes.push(b'\n'); + let result = (|| { + let mut file = open_new_private_file(&temporary)?; + file.write_all(&bytes) + .map_err(|error| StoreError::io(&temporary, error))?; + file.sync_all() + .map_err(|error| StoreError::io(&temporary, error))?; + fs::rename(&temporary, destination).map_err(|error| StoreError::io(destination, error))?; + sync_directory(temporary_parent) + })(); + if result.is_err() { + let _ = fs::remove_file(&temporary); + } + result +} + +pub(crate) fn sync_directory(path: &Path) -> Result<(), StoreError> { + File::open(path) + .and_then(|file| file.sync_all()) + .map_err(|error| StoreError::io(path, error)) +} + +pub(crate) fn random_id() -> Result { + let mut bytes = [0_u8; 16]; + getrandom::fill(&mut bytes).map_err(|error| { + StoreError::new( + StoreErrorKind::Io, + format!("random ID generation failed: {error}"), + ) + })?; + Ok(hex::encode(bytes)) +} + +pub(crate) fn scan_usage(path: &Path) -> Result { + let metadata = fs::symlink_metadata(path).map_err(|error| StoreError::io(path, error))?; + let file_type = metadata.file_type(); + if file_type.is_symlink() { + return Err(StoreError::at( + StoreErrorKind::UnsafeStorePath, + path, + "symlink found while scanning tool-owned storage", + )); + } + if !(file_type.is_file() || file_type.is_dir()) { + return Err(StoreError::at( + StoreErrorKind::CorruptStore, + path, + "special file found while scanning tool-owned storage", + )); + } + let mut usage = StorageUsage { + logical_bytes: if file_type.is_file() { + metadata.len() + } else { + 0 + }, + allocated_bytes: metadata.blocks().saturating_mul(512), + }; + if file_type.is_dir() { + let mut entries = fs::read_dir(path) + .map_err(|error| StoreError::io(path, error))? + .collect::, _>>() + .map_err(|error| StoreError::io(path, error))?; + entries.sort_by_key(std::fs::DirEntry::file_name); + for entry in entries { + let child = scan_usage(&entry.path())?; + usage.logical_bytes = usage + .logical_bytes + .checked_add(child.logical_bytes) + .ok_or_else(|| { + StoreError::new(StoreErrorKind::CorruptStore, "logical usage overflow") + })?; + usage.allocated_bytes = usage + .allocated_bytes + .checked_add(child.allocated_bytes) + .ok_or_else(|| { + StoreError::new(StoreErrorKind::CorruptStore, "allocated usage overflow") + })?; + } + } + Ok(usage) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{ + ffi::OsString, + fs, + os::unix::fs::{MetadataExt as _, PermissionsExt as _, symlink}, + }; + + use tempfile::TempDir; + + use super::AnalysisStore; + use crate::store::{StoreEnvironment, resolve_store}; + + fn store(temp: &TempDir) -> AnalysisStore { + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("initialize") + } + + #[test] + fn initialization_is_private_and_versioned() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + assert_eq!( + fs::metadata(store.root()) + .expect("metadata") + .permissions() + .mode() + & 0o777, + 0o700 + ); + let manifest = fs::read_to_string(store.root().join("store.json")).expect("manifest"); + assert_eq!(manifest, "{\"store_layout_version\":1}\n"); + } + + #[test] + fn usage_counts_logical_and_allocated_bytes() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + fs::write(store.root().join("artifacts/file"), b"abc").expect("write"); + let usage = store.usage().expect("usage"); + assert!(usage.logical_bytes >= 3); + assert!( + usage.allocated_bytes >= fs::metadata(store.root()).expect("metadata").blocks() * 512 + ); + } + + #[test] + fn store_scan_never_follows_symlinks() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let outside = temp.path().join("outside"); + fs::write(&outside, vec![0_u8; 1024]).expect("outside"); + symlink(&outside, store.root().join("artifacts/escape")).expect("symlink"); + assert!(store.usage().is_err()); + } +} diff --git a/src/store/lock.rs b/src/store/lock.rs new file mode 100644 index 0000000..c08f0c6 --- /dev/null +++ b/src/store/lock.rs @@ -0,0 +1,112 @@ +use std::{ + fs::File, + path::{Path, PathBuf}, +}; + +use rustix::fs::{FlockOperation, Mode, OFlags}; + +use crate::domain::Digest; + +use super::{AnalysisStore, StoreError, StoreErrorKind, filesystem::create_private_directory}; + +/// Held nonblocking exclusive per-Analysis OS lock. +#[derive(Debug)] +pub struct AnalysisLock { + _file: File, + path: PathBuf, + sample: Digest, + profile: Digest, +} + +impl AnalysisLock { + /// Immediately acquires `/.lock`, or reports `AnalysisBusy`. + pub fn acquire( + store: &AnalysisStore, + sample: &Digest, + profile: &Digest, + ) -> Result { + let directory = store.path("locks").join(sample.as_str()); + create_private_directory(&directory)?; + let path = directory.join(format!("{}.lock", profile.as_str())); + let descriptor = rustix::fs::open( + &path, + OFlags::RDWR | OFlags::CREATE | OFlags::NOFOLLOW | OFlags::CLOEXEC, + Mode::RUSR | Mode::WUSR, + ) + .map_err(|error| { + StoreError::at( + StoreErrorKind::UnsafeStorePath, + &path, + format!("failed to safely open Analysis lock: {error}"), + ) + })?; + let file = File::from(descriptor); + match rustix::fs::flock(&file, FlockOperation::NonBlockingLockExclusive) { + Ok(()) => Ok(Self { + _file: file, + path, + sample: sample.clone(), + profile: profile.clone(), + }), + Err(error) if error == rustix::io::Errno::WOULDBLOCK => Err(StoreError::at( + StoreErrorKind::AnalysisBusy, + &path, + format!("Analysis {} is locked by another process", profile.as_str()), + )), + Err(error) => Err(StoreError::at( + StoreErrorKind::Io, + &path, + format!("failed to lock Analysis: {error}"), + )), + } + } + + /// Returns the persistent lock-file path. + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + /// Returns the locked Analysis Profile digest. + #[must_use] + pub const fn profile(&self) -> &Digest { + &self.profile + } + + /// Returns whether this token guards the exact requested Analysis identity. + #[must_use] + pub fn guards(&self, sample: &Digest, profile: &Digest) -> bool { + &self.sample == sample && &self.profile == profile + } +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{ffi::OsString, str::FromStr as _}; + + use tempfile::TempDir; + + use super::AnalysisLock; + use crate::{ + domain::Digest, + store::{AnalysisStore, StoreEnvironment, resolve_store}, + }; + + #[test] + fn second_process_style_lock_attempt_fails_immediately() { + let temp = TempDir::new().expect("temp"); + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from(temp.path().join("store"))), + ..StoreEnvironment::default() + }; + let store = AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store"); + let sample = Digest::from_str(&"a".repeat(64)).expect("digest"); + let profile = Digest::from_str(&"b".repeat(64)).expect("digest"); + let first = AnalysisLock::acquire(&store, &sample, &profile).expect("first"); + assert!(AnalysisLock::acquire(&store, &sample, &profile).is_err()); + drop(first); + assert!(AnalysisLock::acquire(&store, &sample, &profile).is_ok()); + } +} diff --git a/src/store/mod.rs b/src/store/mod.rs new file mode 100644 index 0000000..66e6544 --- /dev/null +++ b/src/store/mod.rs @@ -0,0 +1,28 @@ +//! Secure, synchronous Sample snapshot and Analysis Store primitives. + +mod analysis; +mod artifact; +mod cleanup; +mod cleanup_command; +mod error; +mod filesystem; +mod lock; +mod resolution; +mod sample; + +pub use analysis::{ + AnalysisManifest, AnalysisProfileManifest, AnalysisProject, CreatedBy, InventoryFile, + ManifestSample, QuarantineOutcome, RebuildDecision, analysis_profile_digest, inventory_project, + rebuild_decision, validate_analysis, write_analysis_manifest, +}; +pub use artifact::{Artifact, publish_artifact}; +pub use cleanup::{CleanupPlan, CleanupScope, CleanupTransaction, plan_cleanup}; +pub use cleanup_command::execute_cleanup; +pub use error::{StoreError, StoreErrorKind}; +pub use filesystem::{AnalysisStore, STORE_LAYOUT_VERSION, StoreUsage}; +pub use lock::AnalysisLock; +pub use resolution::{ResolvedStore, StoreEnvironment, resolve_store}; +pub use sample::{DEFAULT_STORE_RESERVE_BYTES, SampleSnapshot, SnapshotOptions, snapshot_sample}; + +pub(crate) const DIRECTORY_MODE: u32 = 0o700; +pub(crate) const FILE_MODE: u32 = 0o600; diff --git a/src/store/resolution.rs b/src/store/resolution.rs new file mode 100644 index 0000000..726a75c --- /dev/null +++ b/src/store/resolution.rs @@ -0,0 +1,166 @@ +use std::{ + ffi::OsString, + path::{Component, Path, PathBuf}, +}; + +use crate::{ + domain::TaggedPath, + operation::{AnalysisStoreProvenance, StoreSource}, +}; + +use super::{StoreError, StoreErrorKind}; + +/// Environment values that participate in deterministic store resolution. +#[derive(Clone, Debug, Default)] +pub struct StoreEnvironment { + /// Exact `GHIDR_STORE` value, when set. + pub ghidr_store: Option, + /// Exact `XDG_CACHE_HOME` value, when explicitly set. + pub xdg_cache_home: Option, + /// Exact home directory used only for the final fallback. + pub home: Option, +} + +impl StoreEnvironment { + /// Captures the three relevant variables without lossy decoding. + #[must_use] + pub fn from_process() -> Self { + Self { + ghidr_store: std::env::var_os("GHIDR_STORE"), + xdg_cache_home: std::env::var_os("XDG_CACHE_HOME"), + home: std::env::var_os("HOME"), + } + } +} + +/// A validated absolute UTF-8 store path and its precedence source. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ResolvedStore { + path: PathBuf, + source: StoreSource, +} + +impl ResolvedStore { + /// Returns the absolute UTF-8 store path. + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + /// Returns the selected precedence source. + #[must_use] + pub const fn source(&self) -> StoreSource { + self.source + } + + /// Returns the public exact-byte path representation. + #[must_use] + pub fn tagged_path(&self) -> TaggedPath { + TaggedPath::from_path(&self.path) + } + + /// Builds the fixed public provenance object for this resolution. + #[must_use] + pub fn provenance(&self) -> AnalysisStoreProvenance { + AnalysisStoreProvenance { + path: self.tagged_path(), + source: self.source, + } + } +} + +/// Resolves the store using CLI, `GHIDR_STORE`, explicit XDG, then home precedence. +pub fn resolve_store( + cli_path: Option<&Path>, + environment: &StoreEnvironment, +) -> Result { + let (path, source) = if let Some(path) = cli_path { + (path.to_path_buf(), StoreSource::Cli) + } else if let Some(path) = &environment.ghidr_store { + (PathBuf::from(path), StoreSource::Environment) + } else if let Some(path) = &environment.xdg_cache_home { + (PathBuf::from(path).join("ghidr"), StoreSource::Xdg) + } else if let Some(path) = &environment.home { + ( + PathBuf::from(path).join("ghidr-store"), + StoreSource::HomeFallback, + ) + } else { + return Err(StoreError::new( + StoreErrorKind::InvalidStorePath, + "no Analysis Store path is available: HOME is not set", + )); + }; + validate_store_path(&path)?; + Ok(ResolvedStore { path, source }) +} + +fn validate_store_path(path: &Path) -> Result<(), StoreError> { + let Some(text) = path.to_str() else { + return Err(StoreError::at( + StoreErrorKind::InvalidStorePath, + path, + "Analysis Store path must be valid UTF-8", + )); + }; + if !path.is_absolute() + || text.contains(['\n', '\r']) + || path + .components() + .any(|component| matches!(component, Component::CurDir | Component::ParentDir)) + { + return Err(StoreError::at( + StoreErrorKind::InvalidStorePath, + path, + "Analysis Store path must be absolute, normalized UTF-8", + )); + } + Ok(()) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{ffi::OsString, path::Path}; + + use super::{StoreEnvironment, resolve_store}; + use crate::operation::StoreSource; + + #[test] + fn resolution_obeys_fixed_precedence() { + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from("/environment")), + xdg_cache_home: Some(OsString::from("/xdg")), + home: Some(OsString::from("/home/user")), + }; + let resolved = resolve_store(Some(Path::new("/cli")), &environment).expect("valid"); + assert_eq!(resolved.path(), Path::new("/cli")); + assert_eq!(resolved.source(), StoreSource::Cli); + + let resolved = resolve_store(None, &environment).expect("valid"); + assert_eq!(resolved.path(), Path::new("/environment")); + } + + #[cfg(unix)] + #[test] + fn rejects_exact_non_utf8_configured_path() { + use std::os::unix::ffi::OsStringExt as _; + + let environment = StoreEnvironment { + ghidr_store: Some(OsString::from_vec(vec![b'/', b'x', 0xff])), + ..StoreEnvironment::default() + }; + let error = resolve_store(None, &environment).expect_err("must reject"); + assert!(error.path().is_some()); + } + + #[test] + fn invalid_explicit_xdg_does_not_fall_back() { + let environment = StoreEnvironment { + xdg_cache_home: Some(OsString::from("relative")), + home: Some(OsString::from("/home/user")), + ..StoreEnvironment::default() + }; + assert!(resolve_store(None, &environment).is_err()); + } +} diff --git a/src/store/sample.rs b/src/store/sample.rs new file mode 100644 index 0000000..4e41540 --- /dev/null +++ b/src/store/sample.rs @@ -0,0 +1,380 @@ +use std::{ + fs::{self, File}, + io::{Read as _, Write as _}, + os::unix::fs::{FileTypeExt as _, MetadataExt as _}, + path::{Path, PathBuf}, + str::FromStr as _, +}; + +use sha2::{Digest as _, Sha256}; + +use crate::domain::{Digest, TaggedPath}; + +use super::{ + AnalysisStore, StoreError, StoreErrorKind, + filesystem::{create_private_directory, open_new_private_file, random_id, set_file_mode}, +}; + +/// Default free-space reserve retained beyond the observed Sample size. +pub const DEFAULT_STORE_RESERVE_BYTES: u64 = 1_073_741_824; + +/// Bounds governing one staged Sample copy. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct SnapshotOptions { + /// Maximum bytes accepted before and during streaming. + pub max_sample_bytes: u64, + /// Bytes that must remain available beyond the observed Sample size. + pub reserve_bytes: u64, +} + +impl SnapshotOptions { + /// Constructs the production reserve policy for a positive byte ceiling. + pub fn new(max_sample_bytes: u64) -> Result { + if max_sample_bytes == 0 { + return Err(StoreError::new( + StoreErrorKind::SampleTooLarge, + "maximum Sample size must be positive", + )); + } + Ok(Self { + max_sample_bytes, + reserve_bytes: DEFAULT_STORE_RESERVE_BYTES, + }) + } +} + +/// Private read-only Sample snapshot and its content identity. +#[derive(Debug)] +pub struct SampleSnapshot { + digest: Digest, + size_bytes: u64, + path: PathBuf, + staging_directory: PathBuf, + source_path: TaggedPath, +} + +impl SampleSnapshot { + /// SHA-256 of the exact staged bytes. + #[must_use] + pub const fn digest(&self) -> &Digest { + &self.digest + } + + /// Exact staged byte count. + #[must_use] + pub const fn size_bytes(&self) -> u64 { + self.size_bytes + } + + /// Tool-generated UTF-8 path that may be passed to the worker. + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + /// Exact caller-supplied path retained only as provenance. + #[must_use] + pub const fn source_path(&self) -> &TaggedPath { + &self.source_path + } +} + +impl Drop for SampleSnapshot { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.staging_directory); + } +} + +/// Opens, bounds, streams, hashes, verifies, and privately stages one Sample. +pub fn snapshot_sample( + store: &AnalysisStore, + source: &Path, + options: SnapshotOptions, +) -> Result { + snapshot_with_hook(store, source, options, |_| {}) +} + +fn snapshot_with_hook( + store: &AnalysisStore, + source: &Path, + options: SnapshotOptions, + mut after_chunk: impl FnMut(u64), +) -> Result { + if options.max_sample_bytes == 0 { + return Err(StoreError::new( + StoreErrorKind::SampleTooLarge, + "maximum Sample size must be positive", + )); + } + let mut input = File::open(source).map_err(|error| sample_open_error(source, error))?; + let before = input + .metadata() + .map_err(|error| sample_read_error(source, error))?; + if !before.file_type().is_file() + || before.file_type().is_dir() + || before.file_type().is_fifo() + || before.file_type().is_socket() + || before.file_type().is_block_device() + || before.file_type().is_char_device() + { + return Err(StoreError::at( + StoreErrorKind::InvalidSampleType, + source, + "opened Sample is not a regular file", + )); + } + if before.len() > options.max_sample_bytes { + return Err(StoreError::at( + StoreErrorKind::SampleTooLarge, + source, + "Sample exceeds the configured byte limit", + )); + } + require_free_space(store.root(), before.len(), options.reserve_bytes)?; + + let staging_directory = store.path("staging").join(random_id()?); + create_private_directory(&staging_directory)?; + let snapshot_path = staging_directory.join("sample"); + let result = (|| { + let mut output = open_new_private_file(&snapshot_path)?; + let mut hasher = Sha256::new(); + let mut size_bytes = 0_u64; + let mut buffer = vec![0_u8; 1024 * 1024]; + loop { + let count = input + .read(&mut buffer) + .map_err(|error| sample_read_error(source, error))?; + if count == 0 { + break; + } + size_bytes = size_bytes + .checked_add(u64::try_from(count).map_err(|_| { + StoreError::new(StoreErrorKind::SampleTooLarge, "Sample byte count overflow") + })?) + .ok_or_else(|| { + StoreError::new(StoreErrorKind::SampleTooLarge, "Sample byte count overflow") + })?; + if size_bytes > options.max_sample_bytes { + return Err(StoreError::at( + StoreErrorKind::SampleTooLarge, + source, + "Sample grew beyond the configured byte limit while streaming", + )); + } + output + .write_all(&buffer[..count]) + .map_err(|error| StoreError::io(&snapshot_path, error))?; + hasher.update(&buffer[..count]); + after_chunk(size_bytes); + } + output + .sync_all() + .map_err(|error| StoreError::io(&snapshot_path, error))?; + set_file_mode(&snapshot_path, 0o400)?; + let after = input + .metadata() + .map_err(|error| sample_read_error(source, error))?; + if source_metadata_changed(&before, &after) || after.len() != size_bytes { + return Err(StoreError::at( + StoreErrorKind::SampleChanged, + source, + "Sample metadata changed while creating the private snapshot", + )); + } + let digest = Digest::from_str(&hex::encode(hasher.finalize())).map_err(|error| { + StoreError::new( + StoreErrorKind::Io, + format!("computed digest was invalid: {error}"), + ) + })?; + Ok(SampleSnapshot { + digest, + size_bytes, + path: snapshot_path.clone(), + staging_directory: staging_directory.clone(), + source_path: TaggedPath::from_path(source), + }) + })(); + if result.is_err() { + let _ = fs::remove_dir_all(&staging_directory); + } + result +} + +fn require_free_space(path: &Path, observed: u64, reserve: u64) -> Result<(), StoreError> { + let required = observed.checked_add(reserve).ok_or_else(|| { + StoreError::new( + StoreErrorKind::InsufficientStoreSpace, + "required free-space bound overflow", + ) + })?; + let stats = rustix::fs::statvfs(path).map_err(|error| { + StoreError::new( + StoreErrorKind::Io, + format!("statvfs failed for {}: {error}", path.display()), + ) + })?; + let available = stats.f_bavail.saturating_mul(stats.f_frsize); + if available < required { + return Err(StoreError::at( + StoreErrorKind::InsufficientStoreSpace, + path, + format!("Analysis Store has {available} bytes available but {required} are required"), + )); + } + Ok(()) +} + +fn source_metadata_changed(before: &fs::Metadata, after: &fs::Metadata) -> bool { + before.dev() != after.dev() + || before.ino() != after.ino() + || before.len() != after.len() + || before.mtime() != after.mtime() + || before.mtime_nsec() != after.mtime_nsec() + || before.ctime() != after.ctime() + || before.ctime_nsec() != after.ctime_nsec() +} + +fn sample_open_error(path: &Path, error: std::io::Error) -> StoreError { + let kind = if error.kind() == std::io::ErrorKind::NotFound { + StoreErrorKind::SampleNotFound + } else { + StoreErrorKind::SampleUnreadable + }; + StoreError::at(kind, path, format!("failed to open Sample: {error}")) +} + +fn sample_read_error(path: &Path, error: std::io::Error) -> StoreError { + StoreError::at( + StoreErrorKind::SampleUnreadable, + path, + format!("failed to read Sample: {error}"), + ) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod tests { + use std::{ + ffi::OsString, fs, os::unix::ffi::OsStringExt as _, os::unix::fs::PermissionsExt as _, + }; + + use tempfile::TempDir; + + use super::{SnapshotOptions, snapshot_sample, snapshot_with_hook}; + use crate::{ + domain::TaggedPath, + store::{AnalysisStore, StoreEnvironment, resolve_store}, + }; + + fn store(temp: &TempDir) -> AnalysisStore { + let environment = StoreEnvironment { + ghidr_store: Some(temp.path().join("store").into_os_string()), + ..StoreEnvironment::default() + }; + AnalysisStore::initialize(resolve_store(None, &environment).expect("resolve")) + .expect("store") + } + + fn options() -> SnapshotOptions { + SnapshotOptions { + max_sample_bytes: 1024 * 1024, + reserve_bytes: 0, + } + } + + #[test] + fn identical_bytes_have_identical_identity_and_read_only_snapshots() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let first = temp.path().join("first"); + let second = temp.path().join("second"); + fs::write(&first, b"same bytes").expect("first"); + fs::write(&second, b"same bytes").expect("second"); + let first_snapshot = snapshot_sample(&store, &first, options()).expect("snapshot"); + let second_snapshot = snapshot_sample(&store, &second, options()).expect("snapshot"); + assert_eq!(first_snapshot.digest(), second_snapshot.digest()); + assert_eq!( + fs::metadata(first_snapshot.path()) + .expect("metadata") + .permissions() + .mode() + & 0o777, + 0o400 + ); + } + + #[test] + fn streaming_bound_rejects_oversized_sample() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let sample = temp.path().join("large"); + fs::write(&sample, b"12345").expect("write"); + assert!( + snapshot_sample( + &store, + &sample, + SnapshotOptions { + max_sample_bytes: 4, + reserve_bytes: 0 + } + ) + .is_err() + ); + } + + #[test] + fn changed_source_is_rejected_and_staging_is_removed() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let sample = temp.path().join("changing"); + fs::write(&sample, vec![b'a'; 16]).expect("write"); + let mut changed = false; + let result = snapshot_with_hook(&store, &sample, options(), |_| { + if !changed { + fs::write(&sample, vec![b'b'; 17]).expect("change"); + changed = true; + } + }); + assert!(result.is_err()); + assert_eq!( + fs::read_dir(store.root().join("staging")) + .expect("staging") + .count(), + 0 + ); + } + + #[test] + fn non_utf8_source_is_preserved_exactly_as_provenance() { + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let sample = temp.path().join(OsString::from_vec(vec![b's', 0xff])); + fs::write(&sample, b"bytes").expect("write"); + let snapshot = snapshot_sample(&store, &sample, options()).expect("snapshot"); + assert!(matches!( + snapshot.source_path(), + TaggedPath::UnixBytesBase64(_) + )); + assert_eq!( + snapshot.source_path().to_path_buf().expect("decode"), + sample + ); + } + + #[test] + fn replaced_staging_directory_symlink_is_never_traversed() { + use std::os::unix::fs::symlink; + + let temp = TempDir::new().expect("temp"); + let store = store(&temp); + let outside = temp.path().join("outside"); + fs::create_dir(&outside).expect("outside"); + fs::remove_dir(store.root().join("staging")).expect("remove staging"); + symlink(&outside, store.root().join("staging")).expect("attack"); + let sample = temp.path().join("sample"); + fs::write(&sample, b"bytes").expect("sample"); + assert!(snapshot_sample(&store, &sample, options()).is_err()); + assert_eq!(fs::read_dir(outside).expect("outside entries").count(), 0); + } +}