use super::*;
use mj_core::hex::lower_hex;
use std::collections::{BTreeMap, BTreeSet};
pub(in crate::controller) fn replace_installed_worker_binary(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
) -> Result<()> {
let _owner = crate::worker_lifecycle::require(session_id)?;
let plan = installed_worker_binary_replacement_plan(locator, session_id, worker_binary)?;
for command in plan.commands {
execute_checked(executor, command)?;
}
Ok(())
}
pub(crate) struct PreparedWorkerBinary<'a, E: CommandExecutor> {
name: String,
cleanup: CommandSpec,
executor: &'a E,
locator: targets::TargetLocator,
session_id: String,
_live: LiveStaging,
}
struct LiveStaging {
key: (PathBuf, String),
name: String,
}
type LiveStagingRegistry = std::sync::Mutex<BTreeMap<(PathBuf, String), BTreeSet<String>>>;
fn live_stagings() -> &'static LiveStagingRegistry {
static LIVE: std::sync::OnceLock<LiveStagingRegistry> = std::sync::OnceLock::new();
LIVE.get_or_init(Default::default)
}
impl LiveStaging {
fn register(session_id: &str, name: &str) -> (Self, Vec<String>) {
let key = (mj_core::config::data_dir(), session_id.to_owned());
let mut live = live_stagings()
.lock()
.unwrap_or_else(|error| error.into_inner());
let names = live.entry(key.clone()).or_default();
names.insert(name.to_owned());
let names = names.iter().cloned().collect();
(
Self {
key,
name: name.to_owned(),
},
names,
)
}
}
impl Drop for LiveStaging {
fn drop(&mut self) {
let mut live = live_stagings()
.lock()
.unwrap_or_else(|error| error.into_inner());
if let Some(names) = live.get_mut(&self.key) {
names.remove(&self.name);
if names.is_empty() {
live.remove(&self.key);
}
}
}
}
fn sweep_stale_stagings_command(
locator: &targets::TargetLocator,
worker_root: &str,
live: &[String],
) -> CommandSpec {
let script = r#"root=$1; shift
for path in "$root"/hel.prepared-*; do
[ -e "$path" ] || continue
name=${path##*/}
keep=0
for live in "$@"; do
case $name in "$live"|"$live".next) keep=1 ;; esac
done
[ "$keep" = 1 ] || rm -f -- "$path"
done"#;
let mut args = vec![
"sh".to_owned(),
"-c".to_owned(),
script.to_owned(),
"mj-staging-sweep".to_owned(),
worker_root.to_owned(),
];
args.extend(live.iter().cloned());
targets::locator_command(locator, args).purpose("remove stale private worker staging")
}
impl<E: CommandExecutor> PreparedWorkerBinary<'_, E> {
pub(crate) fn install(&self, owner: &crate::worker_lifecycle::WorkerPermit) -> Result<()> {
install_staged_worker_binary(
owner,
&self.name,
self.executor,
&self.locator,
&self.session_id,
)
}
}
pub(crate) fn prepare_recovery_worker_binary<'a, E: CommandExecutor>(
executor: &'a E,
refresh: &DeferredWorkerBinaryRefresh,
) -> Result<Option<PreparedWorkerBinary<'a, E>>> {
let source = worker_binary_for(&refresh.locator, executor)?;
let expected = mj_core::worker_launch::worker_executable_digest(&source)?;
if crate::session_manager::installed_digest_matches(
executor,
&refresh.installed_digest,
&expected,
) {
return Ok(None);
}
stage_worker_binary_for_upgrade(executor, &refresh.locator, &refresh.session_id, &source)
.map(Some)
}
impl<E: CommandExecutor> std::ops::Deref for PreparedWorkerBinary<'_, E> {
type Target = str;
fn deref(&self) -> &str {
&self.name
}
}
impl<E: CommandExecutor> Drop for PreparedWorkerBinary<'_, E> {
fn drop(&mut self) {
if let Err(error) = execute_checked(self.executor, self.cleanup.clone()) {
tracing::warn!(%error, staging = %self.name, "worker binary staging cleanup failed");
}
}
}
pub(in crate::controller) fn file_size(path: &Path) -> Result<u64> {
Ok(std::fs::metadata(path)
.with_context(|| format!("measure {}", path.display()))?
.len())
}
pub(in crate::controller) fn stage_worker_binary_for_upgrade<'a, E: CommandExecutor>(
executor: &'a E,
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
) -> Result<PreparedWorkerBinary<'a, E>> {
let staging = format!(
"hel.prepared-{}",
crate::session_manager::new_command_id("upgrade-stage")?
);
let plan = worker_binary_replacement_plan(locator, session_id, worker_binary, &staging)?;
let root = targets::worker_root(locator, session_id)?;
crate::target_storage::ensure_room_for(
locator,
[format!("{root}/{staging}.next").as_str()],
|| file_size(worker_binary),
"stage the replacement Mjolnir worker",
)?;
let (live, live_names) = LiveStaging::register(session_id, &staging);
let prepared = PreparedWorkerBinary {
cleanup: targets::locator_command(
locator,
vec![
"rm".into(),
"-f".into(),
"--".into(),
format!("{root}/{staging}"),
format!("{root}/{staging}.next"),
],
)
.purpose("discard private worker staging"),
name: staging,
executor,
locator: locator.clone(),
session_id: session_id.to_owned(),
_live: live,
};
execute_checked(
executor,
sweep_stale_stagings_command(locator, &root, &live_names),
)?;
plan.execute(executor)?;
Ok(prepared)
}
pub(in crate::controller) fn install_staged_worker_binary(
owner: &crate::worker_lifecycle::WorkerPermit,
staging: &str,
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
) -> Result<()> {
ensure!(
owner.session_id() == session_id,
"prepared worker owner mismatch"
);
ensure!(
staging.starts_with("hel.prepared-") && !staging.contains('/'),
"invalid worker staging name"
);
let root = targets::worker_root(locator, session_id)?;
execute_checked(
executor,
targets::locator_command(
locator,
vec![
"mv".into(),
"-f".into(),
"--".into(),
format!("{root}/{staging}"),
format!("{root}/hel"),
],
)
.purpose("install the prepared Mjolnir worker"),
)?;
Ok(())
}
pub(in crate::controller) fn replace_installed_worker_launch_config(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
launch: &WorkerLaunchConfig,
) -> Result<()> {
let _owner = crate::worker_lifecycle::require(session_id)?;
let plan = worker_launch_refresh_plan(locator, session_id, launch)?;
for command in plan.replace.commands {
execute_checked(executor, command)?;
}
Ok(())
}
pub(in crate::controller) fn prepare_managed_harness_for_upgrade(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
launch: &WorkerLaunchConfig,
) -> Result<()> {
if !launch.requires_harness_preparation() {
return Ok(());
}
crate::target_storage::ensure_room_for(
locator,
[
targets::worker_root(locator, session_id)?.as_str(),
crate::target_storage::harness_cache_path(locator).as_str(),
launch.harness_home.to_string_lossy().as_ref(),
],
|| file_size(worker_binary),
"prepare the managed harness",
)?;
if locator.container_engine().is_some() {
return prepare_container_harness_for_upgrade(
executor,
locator,
session_id,
worker_binary,
launch,
);
}
let worker_root = targets::worker_root(locator, session_id)?;
let staging_root = format!(
"{worker_root}/{}",
crate::session_manager::new_command_id("harness-prepare")?
);
let staging_binary = format!("{staging_root}/hel");
let staging_config = format!("{staging_root}/launch.json");
let staging = tempfile::tempdir().context("create managed harness upgrade staging")?;
let local_config = staging.path().join("launch.json");
launch.write(&local_config)?;
if matches!(locator, targets::TargetLocator::LocalBare { .. }) {
execute_checked(
executor,
CommandSpec::new(
worker_binary.to_string_lossy().into_owned(),
[
"worker".to_owned(),
"prepare-harness".to_owned(),
"--config".to_owned(),
local_config.to_string_lossy().into_owned(),
],
)
.purpose("prepare exact managed harness"),
)?;
return Ok(());
}
let ssh = match locator {
targets::TargetLocator::AwsEc2 { ssh, .. }
| targets::TargetLocator::SshBare { ssh, .. } => ssh,
_ => bail!("managed harness policy requires a local bare, SSH-bare, or EC2 target"),
};
let result = (|| -> Result<()> {
execute_checked(
executor,
crate::targets::ssh_command(ssh, ["rm", "-rf", "--", &staging_root])
.purpose("clear managed harness preparation staging"),
)?;
execute_checked(
executor,
crate::targets::ssh_command(ssh, ["mkdir", "-p", &staging_root])
.purpose("create managed harness preparation staging"),
)?;
execute_checked(
executor,
crate::targets::scp_upload(ssh, worker_binary, &staging_binary, false)
.purpose("stage current worker for managed harness preparation"),
)?;
execute_checked(
executor,
crate::targets::scp_upload(ssh, &local_config, &staging_config, false)
.purpose("stage managed harness launch configuration"),
)?;
execute_checked(
executor,
crate::targets::ssh_command(ssh, ["chmod", "700", &staging_binary])
.purpose("make managed harness preparation worker executable"),
)?;
execute_checked(
executor,
crate::targets::ssh_command(
ssh,
[
staging_binary.as_str(),
"worker",
"prepare-harness",
"--config",
staging_config.as_str(),
],
)
.purpose("prepare exact managed harness"),
)?;
Ok(())
})();
result.with_context(|| {
format!("harness preparation failed; staging retained at {staging_root}")
})?;
execute_checked(
executor,
crate::targets::ssh_command(ssh, ["rm", "-rf", "--", &staging_root])
.purpose("remove managed harness preparation staging"),
)
.context("clean managed harness preparation staging")?;
Ok(())
}
fn prepare_container_harness_for_upgrade(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
launch: &WorkerLaunchConfig,
) -> Result<()> {
let staging = tempfile::Builder::new()
.prefix("harness-prepare-")
.tempdir()?;
let name = staging
.path()
.file_name()
.context("harness staging directory has no name")?
.to_string_lossy();
let worker_root = targets::worker_root(locator, session_id)?;
let staging_root = format!("{worker_root}/{name}");
let staging_binary = format!("{staging_root}/hel");
let staging_config = format!("{staging_root}/launch.json");
let upload =
worker_binary_replacement_plan(locator, session_id, worker_binary, &format!("{name}/hel"))?;
let result = (|| -> Result<()> {
execute_checked(
executor,
targets::locator_command(
locator,
vec!["mkdir".into(), "-p".into(), staging_root.clone()],
)
.purpose("create container harness preparation staging"),
)?;
upload.execute(executor)?;
execute_checked(
executor,
write_launch_config_command(
locator,
&staging_config,
serde_json::to_vec_pretty(launch)?,
)
.purpose("stage container harness launch configuration"),
)?;
execute_checked(
executor,
targets::locator_command(
locator,
vec![
staging_binary,
"worker".into(),
"prepare-harness".into(),
"--config".into(),
staging_config,
],
)
.purpose("prepare exact container harness before worker upgrade"),
)?;
Ok(())
})();
result.with_context(|| {
format!("container harness preparation failed; staging retained at {staging_root}")
})?;
execute_checked(
executor,
targets::locator_command(
locator,
vec!["rm".into(), "-rf".into(), "--".into(), staging_root],
)
.purpose("remove container harness preparation staging"),
)
.context("clean container harness preparation staging")?;
Ok(())
}
pub(super) fn prepare_installed_managed_harness(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
worker_root: &str,
launch: &WorkerLaunchConfig,
) -> Result<()> {
if !launch.requires_harness_preparation() {
return Ok(());
}
crate::target_storage::ensure_room_for(
locator,
[
crate::target_storage::harness_cache_path(locator).as_str(),
launch.harness_home.to_string_lossy().as_ref(),
],
|| Ok(0),
"install the managed harness",
)?;
let worker_binary = format!("{worker_root}/hel");
let launch_config = format!("{worker_root}/launch.json");
let command = targets::locator_command(
locator,
vec![
worker_binary,
"worker".into(),
"prepare-harness".into(),
"--config".into(),
launch_config,
],
);
execute_checked(
executor,
command.purpose("prepare exact managed harness before worker startup"),
)?;
Ok(())
}
pub(super) fn installed_worker_binary_replacement_plan(
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
) -> Result<CommandPlan> {
worker_binary_replacement_plan(locator, session_id, worker_binary, "hel")
}
fn worker_binary_replacement_plan(
locator: &targets::TargetLocator,
session_id: &str,
worker_binary: &Path,
installed_name: &str,
) -> Result<CommandPlan> {
verify_worker_build(worker_binary)?;
let worker_root = targets::worker_root(locator, session_id)?;
let installed = format!("{worker_root}/{installed_name}");
let staged = format!("{installed}.next");
let commands = match locator {
targets::TargetLocator::LocalBare { .. } => vec![
CommandSpec::new(
"cp",
[worker_binary.to_string_lossy().into_owned(), staged.clone()],
)
.purpose("stage replacement Mjolnir worker"),
CommandSpec::new("mv", ["-f", &staged, &installed])
.purpose("replace installed Mjolnir worker"),
CommandSpec::new("chmod", ["700", &installed])
.purpose("make replaced Mjolnir worker executable"),
],
targets::TargetLocator::LocalPodman { container_id, .. }
| targets::TargetLocator::LocalDocker { container_id, .. }
| targets::TargetLocator::AppleContainer { container_id, .. } => {
let engine = match locator {
targets::TargetLocator::LocalPodman { .. } => "podman",
targets::TargetLocator::LocalDocker { .. } => "docker",
targets::TargetLocator::AppleContainer { .. } => "container",
_ => unreachable!("matched local container target"),
};
vec![
CommandSpec::new(
engine,
[
"cp".into(),
worker_binary.to_string_lossy().into_owned(),
format!("{container_id}:{staged}"),
],
)
.purpose("stage replacement Mjolnir worker"),
CommandSpec::new(
engine,
container_upload_ownership_args(engine, container_id, &worker_root, &[&staged]),
)
.purpose("match replacement worker to the worker directory owner"),
CommandSpec::new(
engine,
[
"exec".into(),
container_id.clone(),
"mv".into(),
"-f".into(),
staged,
installed.clone(),
],
)
.purpose("replace installed Mjolnir worker"),
CommandSpec::new(
engine,
[
"exec".into(),
container_id.clone(),
"chmod".into(),
"700".into(),
installed,
],
)
.purpose("make replaced Mjolnir worker executable"),
]
}
targets::TargetLocator::AwsEc2 { ssh, .. }
| targets::TargetLocator::SshBare { ssh, .. } => vec![
crate::targets::scp_upload(ssh, worker_binary, &staged, false)
.purpose("stage replacement Mjolnir worker"),
crate::targets::ssh_command(ssh, ["mv", "-f", "--", &staged, &installed])
.purpose("replace installed Mjolnir worker"),
crate::targets::ssh_command(ssh, ["chmod", "700", &installed])
.purpose("make replaced Mjolnir worker executable"),
],
targets::TargetLocator::SshPodman {
ssh, container_id, ..
}
| targets::TargetLocator::SshDocker {
ssh, container_id, ..
} => {
let engine = match locator {
targets::TargetLocator::SshPodman { .. } => "podman",
targets::TargetLocator::SshDocker { .. } => "docker",
_ => unreachable!("matched remote container target"),
};
let upload = format!(
"{}/{}",
targets::REMOTE_UPLOAD_STAGING,
crate::session_manager::new_command_id("worker-upload")?
);
vec![
crate::targets::ssh_command(ssh, ["mkdir", "-p", targets::REMOTE_UPLOAD_STAGING])
.purpose("create remote replacement worker staging"),
crate::targets::scp_upload(ssh, worker_binary, &upload, false)
.purpose("stage replacement Mjolnir worker"),
crate::targets::ssh_command(
ssh,
[engine, "cp", &upload, &format!("{container_id}:{staged}")],
)
.purpose("stage replacement Mjolnir worker"),
crate::targets::ssh_command(
ssh,
std::iter::once(engine.to_owned()).chain(container_upload_ownership_args(
engine,
container_id,
&worker_root,
&[&staged],
)),
)
.purpose("match replacement worker to the worker directory owner"),
crate::targets::ssh_command(
ssh,
[
engine,
"exec",
container_id,
"mv",
"-f",
"--",
&staged,
&installed,
],
)
.purpose("replace installed Mjolnir worker"),
crate::targets::ssh_command(
ssh,
[engine, "exec", container_id, "chmod", "700", &installed],
)
.purpose("make replaced Mjolnir worker executable"),
crate::targets::ssh_command(ssh, ["rm", "-f", "--", &upload])
.purpose("remove remote replacement worker staging"),
]
}
};
Ok(CommandPlan {
description: format!("replace stale Mjolnir worker for session {session_id}"),
commands,
})
}
pub(super) fn installed_file_digest_command(
locator: &targets::TargetLocator,
path: &str,
purpose: &str,
) -> CommandSpec {
let script = "case $(uname -s) in Darwin) exec shasum -a 256 -- \"$1\";; Linux) exec sha256sum -- \"$1\";; *) echo 'unsupported target operating system for worker digest' >&2; exit 1;; esac";
targets::locator_command(
locator,
vec![
"sh".into(),
"-c".into(),
script.into(),
"mj-worker-digest".into(),
path.into(),
],
)
.purpose(purpose)
}
pub(super) fn worker_launch_refresh_plan(
locator: &targets::TargetLocator,
session_id: &str,
launch: &WorkerLaunchConfig,
) -> Result<WorkerLaunchRefreshPlan> {
let worker_root = targets::worker_root(locator, session_id)?;
let installed = format!("{worker_root}/launch.json");
let body = serde_json::to_vec_pretty(launch).context("serialize worker launch config")?;
let expected_sha256 = lower_hex(Sha256::digest(&body));
let replace = write_launch_config_command(locator, &installed, body)
.purpose("replace stale Mjolnir worker launch config");
Ok(WorkerLaunchRefreshPlan {
expected_sha256,
installed_digest: installed_file_digest_command(
locator,
&installed,
"identify installed Mjolnir worker launch config",
),
replace: CommandPlan {
description: format!("replace stale Mjolnir launch config for session {session_id}"),
commands: vec![replace],
},
})
}
fn write_launch_config_command(
locator: &targets::TargetLocator,
installed: &str,
body: Vec<u8>,
) -> CommandSpec {
let staged = format!("{installed}.next");
let staged_arg = targets::join_remote_command(&[staged]);
let installed_arg = targets::join_remote_command(&[installed.to_owned()]);
let script = format!("umask 077; cat > {staged_arg} && mv -f -- {staged_arg} {installed_arg}");
targets::locator_command(locator, vec!["sh".into(), "-c".into(), script])
.with_sensitive_stdin(body)
}
pub(super) fn worker_binary_refresh_plan(
locator: &targets::TargetLocator,
session_id: &str,
) -> Result<Option<WorkerBinaryRefresh>> {
let worker_root = targets::worker_root(locator, session_id)?;
let installed = format!("{worker_root}/hel");
Ok(Some(WorkerBinaryRefresh::Deferred(
DeferredWorkerBinaryRefresh {
locator: locator.clone(),
session_id: session_id.to_owned(),
installed_digest: installed_file_digest_command(
locator,
&installed,
"identify installed Mjolnir worker binary",
),
},
)))
}
pub(crate) fn refresh_target_worker_binary_if_stale(
executor: &impl CommandExecutor,
refresh: &DeferredWorkerBinaryRefresh,
) -> Result<()> {
let source = worker_binary_for(&refresh.locator, executor)
.context("resolve the worker binary for the recovering target")?;
replace_target_worker_binary_if_stale(
executor,
&refresh.locator,
&refresh.session_id,
&refresh.installed_digest,
&source,
)
.map(|_| ())
}
pub(in crate::controller) fn refresh_installed_worker_binary(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
) -> Result<()> {
let Some(WorkerBinaryRefresh::Deferred(refresh)) =
worker_binary_refresh_plan(locator, session_id)?
else {
anyhow::bail!("the worker binary refresh plan was not deferred");
};
refresh_target_worker_binary_if_stale(executor, &refresh)
}
pub(super) fn replace_target_worker_binary_if_stale(
executor: &impl CommandExecutor,
locator: &targets::TargetLocator,
session_id: &str,
installed_digest: &CommandSpec,
source: &Path,
) -> Result<bool> {
verify_worker_build(source)?;
let expected = mj_core::worker_launch::worker_executable_digest(source)?;
let installed = executor
.execute(installed_digest)
.context("read the installed worker digest")?;
let matches = installed.status == 0
&& String::from_utf8_lossy(&installed.stdout)
.split_whitespace()
.next()
.is_some_and(|digest| digest.eq_ignore_ascii_case(&expected));
if matches {
return Ok(false);
}
crate::target_storage::ensure_room_for(
locator,
[targets::worker_root(locator, session_id)?.as_str()],
|| file_size(source),
"replace the stale Mjolnir worker",
)?;
installed_worker_binary_replacement_plan(locator, session_id, source)?
.execute(executor)
.context("replace stale relay worker binary")?;
Ok(true)
}