use super::*;
use mj_core::state::BuildCacheApplication;
use sha2::{Digest, Sha256};
#[cfg(all(test, target_os = "linux"))]
const DEFAULT_MARKER: &str = "# mj automatic shared budget: ";
pub(super) fn shared_directory(cache: &Path) -> PathBuf {
mj_core::config::build_cache_configuration_directory(cache)
}
pub(super) fn host_directory(host: &CacheHost, executor: &impl CommandExecutor) -> Result<PathBuf> {
if matches!(host, CacheHost::Local) {
return dirs::config_dir()
.map(|path| path.join("mbx"))
.context("locate host mbx configuration");
}
let command = host.shell_command(
r#"printf '%s' "${XDG_CONFIG_HOME:-$HOME/.config}/mbx""#,
LABEL,
[],
"locate host mbx configuration",
);
let output = checked(executor.execute(&command)?, &command)?;
let path =
PathBuf::from(String::from_utf8(output.stdout).context("decode mbx configuration path")?);
ensure!(
path.is_absolute(),
"host mbx configuration directory is not absolute"
);
Ok(path)
}
pub(super) fn read_file(
host: &CacheHost,
path: &Path,
executor: &impl CommandExecutor,
) -> Result<Option<String>> {
let command = host.shell_command(
READ_CONFIG_SCRIPT,
LABEL,
[path.to_string_lossy().into_owned()],
"read machine mbx configuration",
);
let output = executor.execute(&command)?;
if output.status == 3 {
return Ok(None);
}
let output = checked(output, &command)?;
Ok(Some(
String::from_utf8(output.stdout).context("decode machine mbx configuration")?,
))
}
#[cfg(all(test, target_os = "linux"))]
pub(super) fn managed_document(settings: &TargetBuildCache, automatic: &str) -> Result<String> {
#[derive(serde::Serialize)]
struct Document<'a> {
gc: Gc<'a>,
#[serde(skip_serializing_if = "Scheduler::is_empty")]
scheduler: Scheduler<'a>,
}
#[derive(serde::Serialize)]
struct Gc<'a> {
max_total_size: &'a str,
}
#[derive(serde::Serialize)]
struct Scheduler<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
cpus: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
memory: Option<&'a str>,
}
impl Scheduler<'_> {
fn is_empty(&self) -> bool {
self.cpus.is_none() && self.memory.is_none()
}
}
let document = Document {
gc: Gc {
max_total_size: settings.max_total_size.as_deref().unwrap_or(automatic),
},
scheduler: Scheduler {
cpus: settings.scheduler.cpus,
memory: settings.scheduler.memory.as_deref(),
},
};
Ok(format!(
"# Managed by Mjolnir. Change the machine's Build cache settings.\n{DEFAULT_MARKER}{automatic}\n{}",
toml::to_string(&document)?
))
}
pub(super) fn configured_limit(
text: Option<&str>,
table: &str,
field: &str,
) -> Result<Option<String>> {
let Some(text) = text else {
return Ok(None);
};
let document: toml::Value = toml::from_str(text).context("parse host mbx configuration")?;
let Some(value) = document.get(table).and_then(|table| table.get(field)) else {
return Ok(None);
};
let value = value
.as_str()
.with_context(|| format!("mbx {table}.{field} must be a size string"))?;
ensure!(
value == "none" || mj_core::config::parse_build_cache_size(value).is_some(),
"mbx {table}.{field} is not a size: {value:?}"
);
Ok(Some(value.to_owned()))
}
const APPLY_SCRIPT: &str = r#"set -eu
directory=$1
expected=$2
mkdir -p -- "$directory"
exec 9>"$directory/.mj-apply.lock"
flock -w 5 9
file="$directory/config.toml"
actual=missing
if [ -f "$file" ]; then actual=$(sha256sum -- "$file"); actual=${actual%% *}; fi
[ "$actual" = "$expected" ] || exit 75
temporary=$(mktemp "$directory/.config.XXXXXX")
trap 'rm -f -- "$temporary"' EXIT HUP INT TERM
cat > "$temporary"
chmod 600 "$temporary"
if [ -f "$file" ] && cmp -s -- "$temporary" "$file"; then exit 0; fi
mv -f -- "$temporary" "$file"
"#;
pub(super) fn apply(
host: &CacheHost,
cache: &ResolvedBuildCache,
executor: &impl CommandExecutor,
) -> Result<bool> {
if cache.previous_config == cache.config_file {
return Ok(true);
}
let expected = cache
.previous_config
.as_deref()
.map(|text| mj_core::hex::lower_hex(Sha256::digest(text.as_bytes())))
.unwrap_or_else(|| "missing".into());
let command = host
.shell_command(
APPLY_SCRIPT,
LABEL,
[
cache.config_directory.to_string_lossy().into_owned(),
expected,
],
"apply machine build cache budgets",
)
.with_sensitive_stdin(
cache
.config_file
.as_deref()
.context("missing managed mbx configuration")?
.as_bytes()
.to_vec(),
);
let output = executor.execute(&command)?;
if output.status == 75 {
return Ok(false);
}
checked(output, &command)?;
Ok(true)
}
pub(super) fn application(previous: Option<&str>, desired: Option<&str>) -> BuildCacheApplication {
if previous == desired {
BuildCacheApplication::Applied
} else {
BuildCacheApplication::Pending
}
}
#[cfg(all(test, target_os = "linux"))]
mod tests {
use super::*;
use std::os::unix::fs::{MetadataExt, symlink};
fn cache(directory: &Path, previous: Option<String>, text: String) -> ResolvedBuildCache {
ResolvedBuildCache {
directory: directory.to_owned(),
native_mbx: NativeMbx {
program: PathBuf::from("/usr/local/bin/mbx"),
version: MBX_VERSION.into(),
},
target_root: None,
config_directory: shared_directory(directory),
config_file: Some(text),
previous_config: previous,
}
}
#[test]
fn both_consumers_see_atomic_budget_updates_without_relinking() {
let temporary = tempfile::tempdir().unwrap();
let executor = targets::ProcessExecutor;
let first = managed_document(
&TargetBuildCache {
max_total_size: Some("500GiB".into()),
..Default::default()
},
"100GB",
)
.unwrap();
let initial = cache(temporary.path(), None, first.clone());
assert!(apply(&CacheHost::Local, &initial, &executor).unwrap());
let file = initial.config_directory.join("config.toml");
let paths = [
temporary.path().join("harness.toml"),
temporary.path().join("reviewer.toml"),
];
for path in &paths {
symlink(&file, path).unwrap();
}
let second = managed_document(
&TargetBuildCache {
max_total_size: Some("100GB".into()),
..Default::default()
},
"100GB",
)
.unwrap();
let second = format!("{second}# {}\n", "x".repeat(256 * 1024));
let updated = cache(temporary.path(), Some(first), second.clone());
assert!(apply(&CacheHost::Local, &updated, &executor).unwrap());
for path in &paths {
let text = std::fs::read_to_string(path).unwrap();
assert_eq!(text, second);
assert_eq!(
configured_limit(Some(&text), "target", "max_size")
.unwrap()
.as_deref(),
None
);
assert_eq!(
configured_limit(Some(&text), "gc", "max_total_size")
.unwrap()
.as_deref(),
Some("100GB")
);
}
let inode = std::fs::metadata(&file).unwrap().ino();
assert!(
apply(
&CacheHost::Local,
&cache(temporary.path(), Some(second.clone()), second),
&executor
)
.unwrap()
);
assert_eq!(
std::fs::metadata(&file).unwrap().ino(),
inode,
"unchanged policy must not be rewritten"
);
}
#[test]
fn existing_cache_mounts_receive_machine_changes_without_changing_the_source() {
let temporary = tempfile::tempdir().unwrap();
let host = CacheHost::Local;
let executor = targets::ProcessExecutor;
let source = temporary.path().join("native.toml");
let current = temporary.path().join("current-cache");
let existing = temporary.path().join("existing-container-cache");
for budget in ["150GiB", "250GiB"] {
let text = format!("[target]\nmax_size = '{budget}'\n");
std::fs::write(&source, &text).unwrap();
let policy = cache(¤t, None, text.clone());
super::super::publish_at(&host, &policy, &existing, &executor).unwrap();
assert_eq!(std::fs::read_to_string(&source).unwrap(), text);
assert_eq!(
std::fs::read_to_string(shared_directory(&existing).join("config.toml")).unwrap(),
text,
);
}
}
#[test]
fn stale_application_cannot_overwrite_a_newer_machine_policy() {
let temporary = tempfile::tempdir().unwrap();
let executor = targets::ProcessExecutor;
let first = cache(
temporary.path(),
None,
"[gc]\nmax_total_size = '500GiB'\n".into(),
);
let stale = cache(
temporary.path(),
None,
"[gc]\nmax_total_size = '1GiB'\n".into(),
);
assert!(apply(&CacheHost::Local, &first, &executor).unwrap());
assert!(!apply(&CacheHost::Local, &stale, &executor).unwrap());
assert_eq!(
std::fs::read_to_string(first.config_directory.join("config.toml")).unwrap(),
first.config_file.unwrap()
);
}
}