use super::*;
pub(crate) fn worker_connect_needs_restart(error: &anyhow::Error) -> bool {
RelayTransportDead::marks(error)
}
pub(super) fn worker_connect_allows_live_restart(error: &anyhow::Error) -> bool {
RelayTransportDead::marks_failed_handshake(error)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum WorkerRecoveryOutcome {
Alive,
Starting,
TargetMissing,
Suppressed,
WorkspaceMissing(PathBuf),
StorageFull {
host: String,
problem: String,
},
RestartedDead,
RestartedUnresponsive,
}
fn storage_forbids_restart(
plan: &WorkerRecoveryPlan,
session_id: Option<&str>,
executor: &impl CommandExecutor,
dead: bool,
) -> Option<WorkerRecoveryOutcome> {
let (host, worker_root) =
crate::target_storage::session_worker_root(&plan.source_target, session_id.unwrap_or(""));
if dead && let Some(command) = &plan.exit_record {
match executor.execute(command) {
Ok(output) => {
if let Some((reason, at)) = crate::controller::recorded_exit_reason(&output.stdout)
&& mj_core::targets::storage::reports_no_space(&reason)
{
crate::target_storage::observe_no_space_at(
&host,
Some(&worker_root),
&format!("the worker stopped: {reason}"),
at.unwrap_or_else(|| chrono::Utc::now().timestamp().max(0) as u64),
);
}
}
Err(error) => tracing::warn!(
%host,
error = format!("{error:#}"),
"could not read the dead worker's exit record"
),
}
}
crate::target_storage::worker_root_problem(&plan.source_target, session_id.unwrap_or(""))
.map(|(host, problem)| WorkerRecoveryOutcome::StorageFull { host, problem })
}
pub(super) fn refresh_worker_binary_if_stale(
executor: &impl CommandExecutor,
refresh: Option<&WorkerBinaryRefresh>,
) -> Result<()> {
match refresh {
None => Ok(()),
Some(WorkerBinaryRefresh::Prepared(plan)) => {
mj_core::worker_build::verify_worker_build(&plan.source)?;
let expected = mj_core::worker_launch::worker_executable_digest(&plan.source)?;
if installed_digest_matches(executor, &plan.installed_digest, &expected) {
return Ok(());
}
plan.replace
.execute(executor)
.context("replace stale relay worker binary")?;
Ok(())
}
Some(WorkerBinaryRefresh::Deferred(refresh)) => {
crate::controller::refresh_target_worker_binary_if_stale(executor, refresh)
}
}
}
pub(crate) fn installed_digest_matches(
executor: &impl CommandExecutor,
command: &CommandSpec,
expected: &str,
) -> bool {
executor.execute(command).as_ref().is_ok_and(|output| {
output.status == 0
&& String::from_utf8_lossy(&output.stdout)
.split_whitespace()
.next()
.is_some_and(|digest| digest.eq_ignore_ascii_case(expected))
})
}
pub(super) fn refresh_worker_launch_if_stale(
executor: &impl CommandExecutor,
plan: Option<&WorkerLaunchRefreshPlan>,
) -> Result<()> {
let Some(plan) = plan else {
return Ok(());
};
if installed_digest_matches(executor, &plan.installed_digest, &plan.expected_sha256) {
return Ok(());
}
plan.replace
.execute(executor)
.context("replace stale relay worker launch config")?;
Ok(())
}
#[cfg(test)]
pub(super) async fn recover_worker(
plan: WorkerRecoveryPlan,
restart_unresponsive: bool,
) -> Result<WorkerRecoveryOutcome> {
recover_worker_for_session(plan, restart_unresponsive, None).await
}
pub(super) async fn recover_worker_for_session(
plan: WorkerRecoveryPlan,
restart_unresponsive: bool,
session_id: Option<String>,
) -> Result<WorkerRecoveryOutcome> {
tokio::task::spawn_blocking(move || {
let executor = CancellableProcessExecutor::with_timeout(WORKER_RESTART_TIMEOUT);
recover_worker_controlled(plan, restart_unresponsive, session_id.as_deref(), &executor)
})
.await
.context("worker recovery task failed")?
}
pub(crate) fn recover_worker_controlled(
mut plan: WorkerRecoveryPlan,
restart_unresponsive: bool,
session_id: Option<&str>,
executor: &impl CommandExecutor,
) -> Result<WorkerRecoveryOutcome> {
let owner = match session_id {
Some(id) => match crate::worker_lifecycle::WorkerPermit::try_acquire(id, "relay recovery")?
{
Some(owner) => Some(owner),
None => return Ok(WorkerRecoveryOutcome::Suppressed),
},
None => {
#[cfg(not(test))]
bail!("worker recovery requires a session identity");
#[cfg(test)]
{
None
}
}
};
let mut work = || -> Result<WorkerRecoveryOutcome> {
if let Some(id) = session_id {
let session = crate::database::read_durable_session_record(id)
.context("read durable session before worker recovery")?;
let eligible = session.as_ref().is_some_and(|session| {
crate::pollers::session_target_is_pollable(session)
&& session.target.as_ref() == Some(&plan.source_target)
});
if !eligible || crate::controller::move_session::move_owns_session(id) {
return Ok(WorkerRecoveryOutcome::Suppressed);
}
}
if let Some(id) = session_id
&& let Some(operation) = crate::database::load_move_operation(id)?
&& operation.source_checkpoint_only
&& operation.destination_target.is_none()
{
plan = crate::controller::Controller::load()?
.worker_recovery_plan(id, Some(&operation))?;
}
if ensure_recovery_target_running(executor, plan.target.as_ref())
.context("restore relay worker target")?
== TargetRecoveryOutcome::Missing
{
return Ok(WorkerRecoveryOutcome::TargetMissing);
}
let output = executor
.execute(&plan.liveness_probe)
.context("probe relay worker liveness")?;
if output.status != 0 {
bail!(
"{} failed with status {}: {}",
plan.liveness_probe.purpose,
output.status,
String::from_utf8_lossy(&output.stderr).trim()
);
}
let restart_pending = session_id
.map(crate::database::load_worker_restart)
.transpose()?
.flatten()
.is_some_and(|intent| intent.target == plan.source_target);
match String::from_utf8_lossy(&output.stdout).trim() {
"alive" if restart_pending => Ok(WorkerRecoveryOutcome::Starting),
"starting" => Ok(WorkerRecoveryOutcome::Starting),
"alive" if !restart_unresponsive => Ok(WorkerRecoveryOutcome::Alive),
"alive" => {
if let Some(full) = storage_forbids_restart(&plan, session_id, executor, false) {
return Ok(full);
}
if let Some(workspace) = plan.workspace.as_ref()
&& !crate::controller::path_exists_on_managed_target(
executor,
&workspace.target,
&workspace.directory,
)?
{
return Ok(WorkerRecoveryOutcome::WorkspaceMissing(
workspace.directory.clone(),
));
}
restart_owned_worker(&plan, session_id, executor, false)?;
Ok(WorkerRecoveryOutcome::RestartedUnresponsive)
}
"dead" => {
if let Some(full) = storage_forbids_restart(&plan, session_id, executor, true) {
return Ok(full);
}
if let Some(workspace) = plan.workspace.as_ref()
&& !crate::controller::path_exists_on_managed_target(
executor,
&workspace.target,
&workspace.directory,
)?
{
return Ok(WorkerRecoveryOutcome::WorkspaceMissing(
workspace.directory.clone(),
));
}
restart_owned_worker(&plan, session_id, executor, true)?;
Ok(WorkerRecoveryOutcome::RestartedDead)
}
output => bail!("worker liveness probe returned unexpected output {output:?}"),
}
};
match owner {
Some(owner) => owner.scope_blocking(work),
None => work(),
}
}
fn restart_owned_worker(
plan: &WorkerRecoveryPlan,
session_id: Option<&str>,
executor: &impl CommandExecutor,
proved_dead: bool,
) -> Result<()> {
let owner = session_id
.map(crate::worker_lifecycle::require)
.transpose()?;
if let Some(owner) = owner.as_ref() {
owner.verify_target(&plan.source_target)?;
if proved_dead {
owner.settle_dead_restart(&plan.source_target)?;
}
}
let deferred = match plan.binary_refresh.as_ref() {
Some(WorkerBinaryRefresh::Deferred(refresh)) => Some(refresh),
_ => None,
};
let staged = deferred
.map(|refresh| crate::controller::prepare_recovery_worker_binary(executor, refresh))
.transpose()?
.flatten();
let _swap = crate::upgrade::activity_unless_draining("worker recovery swap")?;
if let Some(owner) = owner.as_ref() {
owner.begin_restart(&plan.source_target, String::new())?;
crate::database::advance_worker_restart(
owner.session_id(),
owner.operation_id(),
crate::database::WorkerRestartPhase::Swapping,
)?;
}
if let Some(staged) = staged {
staged.install(
owner
.as_ref()
.context("prepared recovery has no worker owner")?,
)?;
} else if deferred.is_none() {
refresh_worker_binary_if_stale(executor, plan.binary_refresh.as_ref())?;
}
refresh_worker_launch_if_stale(executor, plan.launch_refresh.as_ref())?;
plan.restart.execute(executor)?;
if let Some(owner) = owner.as_ref() {
crate::database::advance_worker_restart(
owner.session_id(),
owner.operation_id(),
crate::database::WorkerRestartPhase::AwaitingReadiness,
)?;
}
Ok(())
}