brokk-mj-controller 2.31.0

Daemon-side controller, session manager, and web server for Mjolnir
Documentation
use super::*;

/// Whether a relay failure means the transport to the worker is gone, so
/// restarting that worker is the only recovery left.
///
/// Every failure that proves it is marked with [`RelayTransportDead`] where it
/// is produced, and this decision downcasts for that marker. Message text is
/// never read: a reworded diagnostic must not be able to disable auto-restart.
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),
    /// The target's disk is full. Nothing was restarted: a restart writes to
    /// the target, and a worker started there would stop again at once.
    StorageFull {
        host: String,
        problem: String,
    },
    RestartedDead,
    RestartedUnresponsive,
}

/// Whether a full disk forbids restarting this worker. A dead worker's own
/// exit record is read first: if it stopped for lack of space, the storage
/// owner learns that before anything is written, and the record is read
/// before a restart would replace it.
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)
                {
                    // The worker stops when it cannot write its journal,
                    // which lives in the worker root.
                    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(())
        }
        // Pick the binary for the target's architecture and copy only
        // if it differs. Runs here in the recovery task, never on the UI path.
        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> {
    // Background recovery never waits while its actor's control mailbox is
    // needed by the current lifecycle owner.
    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);
            }
        }
        // A failed Move can leave this actor with a plan from before recovery.
        // Never overwrite the durable checkpoint-only launch with that old plan.
        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()
            );
        }
        // A replacement may still be replaying its journal after the daemon
        // that launched it has exited. A lost handshake is not permission to kill
        // that live process; only a liveness probe proving death permits restart.
        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(),
    }
}

/// Preparation runs without daemon admission. Only the bounded process swap
/// holds admission; its durable intent protects journal replay after handoff.
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(())
}