Skip to main content

mj_controller/controller/
worker_restart.rs

1//! Replacing a session's worker process in place.
2//!
3//! Two things ask for this: a checkpoint whose ACP turn will not finish, and a
4//! session whose worker predates the controller now talking to it. Both stop
5//! the worker, install the binary this controller would provision, start it
6//! and reconnect, so the sequence lives here once and each caller supplies
7//! only what it tells the operator.
8
9use std::time::Duration;
10
11use anyhow::{Context, Result, bail};
12
13use crate::session_manager::{SessionManagerControl, StandaloneSession};
14use crate::targets::{self, CommandExecutor, CommandSpec};
15use mj_core::config::HarnessKind;
16use mj_core::relay::RelayExecutionState;
17
18use super::Controller;
19use super::readiness::{connect_started_worker_with_timeout, wait_for_native_session};
20use super::worker_binary::{
21    install_staged_worker_binary, prepare_managed_harness_for_upgrade,
22    replace_installed_worker_binary, replace_installed_worker_launch_config,
23    stage_worker_binary_for_upgrade, start_worker, start_worker_durably,
24    stop_worker_after_target_recovery, worker_binary_for, worker_probe_diagnosis,
25};
26
27/// How long a restarted worker has to recover its journal, bind `control.sock`
28/// and report an idle ACP session. Journal recovery over a long transcript
29/// runs before the socket exists, so this has to outlast it.
30const WORKER_RESTART_TIMEOUT: Duration = Duration::from_secs(300);
31
32/// How long a quiet session's upgrade waits for its actor and its lease.
33const UPGRADE_LEASE_TIMEOUT: Duration = Duration::from_secs(5);
34
35/// The worker was stopped so it could be replaced, and no worker came back:
36/// the binary swap, start, connect, or ACP readiness after it failed. The
37/// session has no live worker until something restarts one.
38#[derive(Debug)]
39pub struct WorkerRestartLeftNoWorker;
40
41impl WorkerRestartLeftNoWorker {
42    /// Whether a failed operation left the session without a live worker.
43    ///
44    /// The marker is carried by the error, not by its text. Callers wrap
45    /// restart errors in further context, and `anyhow`'s downcast walks those
46    /// layers, so added context does not hide it.
47    #[must_use]
48    pub fn marks(error: &anyhow::Error) -> bool {
49        error.downcast_ref::<Self>().is_some()
50    }
51}
52
53impl std::fmt::Display for WorkerRestartLeftNoWorker {
54    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
55        formatter.write_str("the worker restart left the session without a live worker")
56    }
57}
58
59impl std::error::Error for WorkerRestartLeftNoWorker {}
60
61/// After a restarted worker has answered once, only a dead transport proves
62/// the worker is gone again; any other failure leaves a live worker behind.
63fn mark_if_transport_died(error: anyhow::Error) -> anyhow::Error {
64    if crate::worker_client::RelayTransportDead::marks(&error) {
65        error.context(WorkerRestartLeftNoWorker)
66    } else {
67        error
68    }
69}
70
71/// What one restart tells the operator at each step. The steps are identical;
72/// only the reason differs, and a diagnostic that named the wrong reason would
73/// send someone looking in the wrong place.
74pub(super) struct WorkerRestartMessages {
75    pub stop: &'static str,
76    pub replace: &'static str,
77    pub start: &'static str,
78    pub connect: &'static str,
79    pub project_memory: &'static str,
80    pub native_session: &'static str,
81}
82
83pub(super) struct InstalledWorkerRestart<'a> {
84    pub backend: &'a targets::TargetLocator,
85    pub worker_root: &'a str,
86    pub reconnect: &'a CommandSpec,
87    pub launch: Option<&'a mj_core::worker_launch::WorkerLaunchConfig>,
88    pub prepared: bool,
89    pub messages: &'a WorkerRestartMessages,
90}
91
92/// A wedged ACP turn is being killed so a checkpoint barrier can be admitted.
93pub(super) const RESTART_FOR_CHECKPOINT: WorkerRestartMessages = WorkerRestartMessages {
94    stop: "stop wedged Mjolnir worker before retrying checkpoint",
95    replace: "replace Mjolnir worker binary before retrying checkpoint",
96    start: "start Mjolnir worker after interrupting a wedged ACP turn",
97    connect: "connect to Mjolnir worker after restarting it for checkpoint",
98    project_memory: "project memory will not be synchronized after checkpoint worker restart",
99    native_session: "wait for ACP session after restarting the worker for checkpoint",
100};
101
102/// A quiet session is being moved onto the worker binary this controller
103/// would install.
104const RESTART_FOR_UPGRADE: WorkerRestartMessages = WorkerRestartMessages {
105    stop: "stop the Mjolnir worker before installing the current binary",
106    replace: "install the current Mjolnir worker binary",
107    start: "start Mjolnir worker on the current binary",
108    connect: "connect to Mjolnir worker after upgrading its binary",
109    project_memory: "project memory will not be synchronized after the worker upgrade",
110    native_session: "wait for ACP session after upgrading the worker",
111};
112
113/// What an upgrade attempt found. Nothing here is a failure: a worker that is
114/// already current and a session that started working again are both ordinary.
115#[derive(Debug, Clone, PartialEq, Eq)]
116pub enum WorkerUpgradeOutcome {
117    /// The worker was replaced and the managed session now speaks to one
118    /// running this build.
119    Upgraded { build: String },
120    /// The worker already runs the binary this controller would install.
121    AlreadyCurrent { build: String },
122    /// The session was working when the upgrade reached it. A worker restart
123    /// would have killed that work, so nothing was touched.
124    Deferred,
125}
126
127impl WorkerUpgradeOutcome {
128    /// The build the session's worker runs now, or `None` when the attempt
129    /// stood down without establishing one.
130    #[must_use]
131    pub fn build(&self) -> Option<&str> {
132        match self {
133            Self::Upgraded { build } | Self::AlreadyCurrent { build } => Some(build),
134            Self::Deferred => None,
135        }
136    }
137}
138
139/// Whether a worker that reported `reported` in hello is running `installed`,
140/// the binary this controller would provision.
141///
142/// A worker that reported nothing is not: the field postdates it, so its
143/// binary does too.
144fn worker_runs_installed_build(reported: Option<&str>, installed: &str) -> bool {
145    reported.is_some_and(|reported| reported == installed)
146}
147
148impl Controller {
149    /// Replace a session's worker with the binary this controller would
150    /// install, when the session is quiet and its worker is a different build.
151    ///
152    /// `reported_build` is the digest the worker gave the observer that asked
153    /// for this. It only saves work: a match returns before anything is leased.
154    /// The decision that matters is taken again under the lease, against a
155    /// snapshot read from the worker itself, because a session can start
156    /// working between an observation and this call.
157    pub async fn upgrade_session_worker(
158        &self,
159        session_id: &str,
160        executor: &(impl CommandExecutor + Sync),
161        manager: &SessionManagerControl,
162        reported_build: Option<&str>,
163    ) -> Result<WorkerUpgradeOutcome> {
164        let (backend, worker_root) = self.worker_placement(session_id)?;
165        let reconnect = targets::reconnect_plan(&backend, session_id)?
166            .commands
167            .into_iter()
168            .next()
169            .context("reconnect plan is empty")?;
170        let binary = worker_binary_for(&backend, executor)
171            .context("resolve the worker binary this controller would install")?;
172        let installed = mj_core::worker_launch::worker_executable_digest(&binary)?;
173        if worker_runs_installed_build(reported_build, &installed) {
174            return Ok(WorkerUpgradeOutcome::AlreadyCurrent { build: installed });
175        }
176
177        // Preparation can download a managed harness. Keep the old worker and
178        // its controls available for all of it; reserve only for the swap.
179        let launch = self.current_worker_launch_config(session_id, &backend)?;
180        if let Some(session) = self.state.sessions.get(session_id)
181            && session.build_cache.is_some()
182        {
183            self.prepare_build_cache_links(session, &backend, &launch, executor)
184                .context("prepare shared machine cache configuration before worker upgrade")?;
185        }
186        prepare_managed_harness_for_upgrade(executor, &backend, session_id, &binary, &launch)
187            .context("prepare the current managed harness before replacing the worker")?;
188        let staging = stage_worker_binary_for_upgrade(executor, &backend, session_id, &binary)
189            .context("stage the current worker while the old worker remains available")?;
190        let Some(owner) =
191            crate::worker_lifecycle::WorkerPermit::try_acquire(session_id, "worker upgrade")?
192        else {
193            return Ok(WorkerUpgradeOutcome::Deferred);
194        };
195        owner
196            .scope(async {
197                let target = self.state.sessions[session_id]
198                    .target
199                    .as_ref()
200                    .context("worker upgrade has no target")?;
201                owner.verify_target(target)?;
202                let current = Controller::load()?;
203                let current_launch = current.current_worker_launch_config(session_id, &backend)?;
204                if serde_json::to_value(&current_launch)? != serde_json::to_value(&launch)? {
205                    // A Move can change the harness without changing its root.
206                    // Prepared inputs are discarded; the next attempt prepares
207                    // the launch chosen by the current durable session.
208                    return Ok(WorkerUpgradeOutcome::Deferred);
209                }
210                if crate::database::load_worker_restart(session_id)?.is_some() {
211                    return Ok(WorkerUpgradeOutcome::Deferred);
212                }
213                let handle = manager
214                    .wait_for_session(session_id, UPGRADE_LEASE_TIMEOUT)
215                    .await?;
216                let harness = self.state.sessions[session_id].harness_kind;
217                // Reserve handoff admission before taking the worker's atomic idle
218                // reservation. A draining daemon must not take a worker connection.
219                let Ok(swap) = crate::upgrade::activity_unless_draining("worker swap") else {
220                    return Ok(WorkerUpgradeOutcome::Deferred);
221                };
222                let Some(mut lease) =
223                    super::IdleWorkspaceLease::acquire_for_upgrade(&handle, harness).await?
224                else {
225                    return Ok(WorkerUpgradeOutcome::Deferred);
226                };
227                if !lease.verify_for_upgrade().await? {
228                    return Ok(WorkerUpgradeOutcome::Deferred);
229                }
230                // The accepted swap records its target before touching a process. Once
231                // detached startup succeeds, the next daemon can resume observation.
232                let operation_id = owner.operation_id().to_owned();
233                let intent = crate::database::WorkerRestartIntent {
234                    operation_id: operation_id.clone(),
235                    target: self.state.sessions[session_id]
236                        .target
237                        .clone()
238                        .context("worker restart has no durable target")?,
239                    desired_build: installed.clone(),
240                };
241                {
242                    crate::database::begin_worker_restart(session_id, &intent)?;
243                    install_staged_worker_binary(&owner, &staging, executor, &backend, session_id)
244                        .context("install the prepared worker under its idle reservation")?;
245                    replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
246                        .context(
247                            "install the worker launch configuration under its idle reservation",
248                        )?;
249                    crate::database::advance_worker_restart(
250                        session_id,
251                        &operation_id,
252                        crate::database::WorkerRestartPhase::Swapping,
253                    )?;
254                    stop_worker_after_target_recovery(executor, &backend, session_id, &worker_root)
255                        .context(RESTART_FOR_UPGRADE.stop)?;
256                    start_worker_durably(&owner, target, executor, &backend, &worker_root)
257                        .context(RESTART_FOR_UPGRADE.start)
258                        .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
259                }
260                // The process now owns boot/journal recovery. Waiting for its socket
261                // is resumable, and must not hold daemon replacement for minutes.
262                drop(swap);
263                let mut connection = connect_started_worker_with_timeout(
264                    &reconnect,
265                    session_id,
266                    executor,
267                    &backend,
268                    &worker_root,
269                    WORKER_RESTART_TIMEOUT,
270                )
271                .await
272                .context(RESTART_FOR_UPGRADE.connect)?;
273                anyhow::ensure!(
274                    connection.snapshot().worker_build.as_deref() == Some(&installed),
275                    "replacement worker reported an unexpected build"
276                );
277                anyhow::ensure!(
278                    connection.snapshot().operational.checkpoint_only
279                        == (launch.run_mode
280                            == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
281                    "replacement worker reported an unexpected execution mode"
282                );
283                let project_memory = match self.project_memory_sync_target(session_id) {
284                    Ok(target) => Some(target),
285                    Err(error) => {
286                        tracing::warn!(
287                            session_id,
288                            error = format!("{error:#}"),
289                            "project memory will not be synchronized after worker upgrade"
290                        );
291                        None
292                    }
293                };
294                connection.set_project_memory_target(project_memory);
295                wait_for_native_session(&mut connection, executor, harness).await?;
296                wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT, executor).await?;
297                crate::database::finish_worker_restart(session_id, &operation_id)?;
298                lease.finish_replacement(connection);
299                Ok(WorkerUpgradeOutcome::Upgraded { build: installed })
300            })
301            .await
302    }
303
304    /// Stop the worker, install the binary this controller would provision,
305    /// start it and reconnect to the session it recovers.
306    pub(super) async fn restart_worker_with_installed_binary(
307        &self,
308        session_id: &str,
309        executor: &(impl CommandExecutor + Sync),
310        restart: InstalledWorkerRestart<'_>,
311    ) -> Result<StandaloneSession> {
312        crate::worker_lifecycle::run(
313            session_id,
314            "restart worker with installed binary",
315            executor,
316            async {
317                let owner = crate::worker_lifecycle::require(session_id)?;
318                let target = self
319                    .state
320                    .sessions
321                    .get(session_id)
322                    .and_then(|session| session.target.as_ref());
323                if let Some(target) = target {
324                    owner.verify_target(target)?;
325                    if crate::database::load_worker_restart(session_id)?.is_some() {
326                        let output =
327                            executor.execute(&super::worker_binary::worker_liveness_command(
328                                restart.backend,
329                                restart.worker_root,
330                            ))?;
331                        anyhow::ensure!(
332                            output.status == 0
333                                && String::from_utf8_lossy(&output.stdout).trim() == "dead",
334                            "worker replacement is still booting; reconnect before restarting it"
335                        );
336                        owner.settle_dead_restart(target)?;
337                    }
338                    owner.begin_restart(target, String::new())?;
339                    crate::database::advance_worker_restart(
340                        session_id,
341                        owner.operation_id(),
342                        crate::database::WorkerRestartPhase::Swapping,
343                    )?;
344                }
345                // A failed stop may leave the old worker alive, so it stays outside the
346                // marker the start below applies: only steps after a successful stop
347                // can leave the session with no worker at all.
348                stop_worker_after_target_recovery(
349                    executor,
350                    restart.backend,
351                    session_id,
352                    restart.worker_root,
353                )
354                .context(restart.messages.stop)?;
355                self.start_installed_worker(session_id, executor, restart)
356                    .await
357            },
358        )
359        .await
360    }
361
362    /// The part of a restart after its stop: install the binary unless it is
363    /// `prepared`, start the worker on the existing worker root, connect with
364    /// the long restart timeout, and wait until its harness has loaded its
365    /// native session and gone idle. A parked sub-agent is started again with
366    /// exactly this sequence, since its worker was stopped when it was parked.
367    pub(super) async fn start_installed_worker(
368        &self,
369        session_id: &str,
370        executor: &(impl CommandExecutor + Sync),
371        restart: InstalledWorkerRestart<'_>,
372    ) -> Result<StandaloneSession> {
373        crate::worker_lifecycle::run(session_id, "start installed worker", executor, async {
374            let InstalledWorkerRestart {
375                backend,
376                worker_root,
377                reconnect,
378                launch,
379                prepared,
380                messages,
381            } = restart;
382            let owner = crate::worker_lifecycle::require(session_id)?;
383            let session = self.state.sessions.get(session_id);
384            let harness = session
385                .map(|session| session.harness_kind)
386                .or_else(|| launch.map(|launch| launch.harness))
387                .unwrap_or(HarnessKind::Codex);
388            let observed_updated_at = session.map(|session| session.updated_at.clone());
389            let target = self
390                .state
391                .sessions
392                .get(session_id)
393                .and_then(|session| session.target.as_ref());
394            if let Some(target) = target {
395                owner.verify_target(target)?;
396                if let Some(intent) = crate::database::load_worker_restart(session_id)? {
397                    // The checkpoint-only fallback may retry after a proved dead
398                    // boot. Never stop or overwrite a still-live replacement.
399                    if intent.operation_id != owner.operation_id()
400                        || crate::database::worker_restart_phase(session_id)?
401                            == Some(crate::database::WorkerRestartPhase::AwaitingReadiness)
402                    {
403                        let output = executor.execute(
404                            &super::worker_binary::worker_liveness_command(backend, worker_root),
405                        )?;
406                        anyhow::ensure!(
407                            output.status == 0
408                                && String::from_utf8_lossy(&output.stdout).trim() == "dead",
409                            "worker replacement is still booting"
410                        );
411                        owner.settle_dead_restart(target)?;
412                    }
413                }
414                if crate::database::load_worker_restart(session_id)?.is_none() {
415                    owner.begin_restart(target, String::new())?;
416                    crate::database::advance_worker_restart(
417                        session_id,
418                        owner.operation_id(),
419                        crate::database::WorkerRestartPhase::Swapping,
420                    )?;
421                }
422            }
423            // Everything up to the first successful connection either fails with no
424            // worker running or cannot tell: the marker covers all of it.
425            let mut connection = async {
426                // Copy through hel.next and rename. scp/cp onto a still-mapped hel
427                // fails with ETXTBSY ("dest open ... Failure") even after SIGKILL,
428                // and prepare_worker_files writes that path in place.
429                if !prepared {
430                    let binary = worker_binary_for(backend, executor)?;
431                    replace_installed_worker_binary(executor, backend, session_id, &binary)
432                        .context(messages.replace)?;
433                    if let Some(launch) = launch {
434                        replace_installed_worker_launch_config(
435                            executor, backend, session_id, launch,
436                        )
437                        .context("install the current Mjolnir worker launch configuration")?;
438                    }
439                }
440                let started = match target {
441                    Some(target) => {
442                        start_worker_durably(&owner, target, executor, backend, worker_root)
443                    }
444                    None => start_worker(&owner, executor, backend, worker_root),
445                };
446                started.context(messages.start)?;
447                // Journal recovery runs before the daemon binds control.sock. A long
448                // kimi session can take well over the ordinary 30s startup window.
449                match connect_started_worker_with_timeout(
450                    reconnect,
451                    session_id,
452                    executor,
453                    backend,
454                    worker_root,
455                    WORKER_RESTART_TIMEOUT,
456                )
457                .await
458                {
459                    Ok(connection) => Ok(connection),
460                    Err(error) => {
461                        Err(
462                            worker_probe_diagnosis(executor, backend, worker_root, error)
463                                .context(messages.connect),
464                        )
465                    }
466                }
467            }
468            .await
469            .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
470            let project_memory = match self.project_memory_sync_target(session_id) {
471                Ok(target) => Some(target),
472                Err(error) => {
473                    tracing::warn!(
474                        session_id,
475                        error = format!("{error:#}"),
476                        "{}",
477                        messages.project_memory
478                    );
479                    None
480                }
481            };
482            connection.set_project_memory_target(project_memory);
483            // A worker answered, so a failure from here on only means "no worker"
484            // when the transport to it died again.
485            let readiness = async {
486                let checkpoint_only = connection.sync().await?.operational.checkpoint_only;
487                if let Some(launch) = launch {
488                    anyhow::ensure!(
489                        checkpoint_only
490                            == (launch.run_mode
491                                == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
492                        "restarted worker did not enter the requested execution mode"
493                    );
494                }
495                if checkpoint_only {
496                    return Ok(());
497                }
498                wait_for_native_session(&mut connection, executor, harness)
499                    .await
500                    .context(messages.native_session)?;
501                wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT, executor)
502                    .await
503                    .context("wait for ACP to go idle after worker restart")
504            }
505            .await
506            .map_err(mark_if_transport_died);
507            if let Err(error) = readiness {
508                if let (Some(failure), Some(observed_updated_at)) = (
509                    error.downcast_ref::<super::readiness::HarnessPreparationFailure>(),
510                    observed_updated_at.as_deref(),
511                ) {
512                    self.persist_harness_preparation_failure(
513                        session_id,
514                        &failure.to_string(),
515                        observed_updated_at,
516                    )
517                    .await?;
518                }
519                return Err(error);
520            }
521            if target.is_some() {
522                crate::database::finish_worker_restart(session_id, owner.operation_id())?;
523            }
524            Ok(connection)
525        })
526        .await
527    }
528}
529
530/// Wait until a restarted worker's projection stops moving and reports idle.
531///
532/// Three stable polls, not one: a worker that has just recovered its journal
533/// can report idle between two events it is still applying. "Idle" is the
534/// shared predicate, not the bare execution flag, so a foreground tool or a
535/// turn the projection has not caught up with keeps the restart from being
536/// declared ready underneath it. A synchronized active goal is the one
537/// deliberate exception: the restarted worker is meant to continue it.
538async fn wait_for_idle_projection(
539    relay: &mut StandaloneSession,
540    timeout: Duration,
541    executor: &impl CommandExecutor,
542) -> Result<()> {
543    let deadline = tokio::time::Instant::now() + timeout;
544    let mut last_ordinal = None;
545    let mut stable_polls = 0_u8;
546    loop {
547        anyhow::ensure!(
548            !executor.cancellation_requested(),
549            "worker readiness wait cancelled"
550        );
551        let snapshot = relay.sync().await?;
552        let ordinal = snapshot.operational.latest_ordinal;
553        let goal_active =
554            snapshot.operational.goal.synchronized() && snapshot.operational.goal.active();
555        let idle = snapshot.operational.native_session_is_ready()
556            && (!snapshot.operational.has_work_in_flight() || goal_active);
557        if idle && (goal_active || last_ordinal == Some(ordinal)) {
558            stable_polls = stable_polls.saturating_add(1);
559            if stable_polls >= 3 {
560                return Ok(());
561            }
562        } else {
563            stable_polls = 0;
564        }
565        last_ordinal = Some(ordinal);
566        if snapshot.operational.execution == RelayExecutionState::Closed {
567            bail!("ACP runtime stopped before becoming idle");
568        }
569        if tokio::time::Instant::now() >= deadline {
570            bail!(
571                "ACP runtime did not become idle after worker restart (execution={:?}, ordinal={ordinal})",
572                snapshot.operational.execution
573            );
574        }
575        tokio::time::sleep(Duration::from_millis(200)).await;
576    }
577}
578
579#[cfg(test)]
580mod tests {
581    // Hard-won: b9c3a05: issue #1001 left a restarted session marked running after its relay died
582    #[test]
583    fn a_dead_transport_after_reconnect_marks_the_restart_as_leaving_no_worker() {
584        let died = anyhow::Error::new(crate::worker_client::RelayTransportDead::new(
585            "relay proxy disconnected during attach",
586        ))
587        .context("wait for ACP session after restarting the worker for checkpoint");
588        assert!(super::WorkerRestartLeftNoWorker::marks(
589            &super::mark_if_transport_died(died)
590        ));
591
592        let slow = anyhow::anyhow!("timed out waiting for the ACP session")
593            .context("wait for ACP session after restarting the worker for checkpoint");
594        let slow = super::mark_if_transport_died(slow);
595        assert!(!super::WorkerRestartLeftNoWorker::marks(&slow), "{slow:#}");
596    }
597
598    #[cfg(unix)]
599    use super::*;
600
601    #[cfg(unix)]
602    use std::sync::Mutex;
603
604    #[cfg(unix)]
605    use crate::targets::CommandOutput;
606
607    /// Fails every command after the first, so a restart gets past its stop and
608    /// then loses the worker it was replacing.
609    #[cfg(unix)]
610    struct StopSucceedsThenFails {
611        executed: Mutex<Vec<String>>,
612    }
613
614    #[cfg(unix)]
615    impl CommandExecutor for StopSucceedsThenFails {
616        fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
617            let mut executed = self.executed.lock().expect("executed commands");
618            executed.push(command.program.clone());
619            if executed.len() == 1 {
620                return Ok(CommandOutput {
621                    status: 0,
622                    stdout: Vec::new(),
623                    stderr: Vec::new(),
624                });
625            }
626            Ok(CommandOutput {
627                status: 1,
628                stdout: Vec::new(),
629                stderr: b"no such target".to_vec(),
630            })
631        }
632    }
633
634    #[cfg(unix)]
635    struct FailingStop;
636
637    #[cfg(unix)]
638    impl CommandExecutor for FailingStop {
639        fn execute(&self, _command: &CommandSpec) -> Result<CommandOutput> {
640            Ok(CommandOutput {
641                status: 1,
642                stdout: Vec::new(),
643                stderr: b"permission denied".to_vec(),
644            })
645        }
646    }
647
648    #[cfg(unix)]
649    fn bare_restart_controller() -> Controller {
650        Controller {
651            config: mj_core::config::Config::default(),
652            state: mj_core::state::State::default(),
653        }
654    }
655
656    #[cfg(unix)]
657    async fn restart_error(
658        session_id: &str,
659        executor: &(impl CommandExecutor + Sync),
660    ) -> anyhow::Error {
661        let worker_root = format!("/tmp/mjolnir-restart-test/{session_id}");
662        let backend = targets::TargetLocator::LocalBare {
663            worker_root: worker_root.clone(),
664        };
665        let reconnect = CommandSpec::new("unused", std::iter::empty::<&str>());
666        let result = bare_restart_controller()
667            .restart_worker_with_installed_binary(
668                session_id,
669                executor,
670                InstalledWorkerRestart {
671                    backend: &backend,
672                    worker_root: &worker_root,
673                    reconnect: &reconnect,
674                    launch: None,
675                    prepared: false,
676                    messages: &RESTART_FOR_CHECKPOINT,
677                },
678            )
679            .await;
680        match result {
681            Ok(_) => panic!("a failing executor unexpectedly restarted the worker"),
682            Err(error) => error,
683        }
684    }
685
686    #[cfg(unix)]
687    #[tokio::test]
688    async fn a_restart_that_could_not_stop_the_worker_leaves_it_running() {
689        let error = restart_error("0123456789abcdef0123456789abcdef", &FailingStop).await;
690
691        assert!(
692            !WorkerRestartLeftNoWorker::marks(&error),
693            "a failed stop may leave the old worker alive: {error:#}"
694        );
695    }
696
697    #[cfg(unix)]
698    // Hard-won: b9c3a05: issue #1001 left a session marked running after stop succeeded but restart failed
699    #[tokio::test]
700    async fn a_restart_that_stopped_the_worker_and_then_failed_is_marked() {
701        let executor = StopSucceedsThenFails {
702            executed: Mutex::new(Vec::new()),
703        };
704
705        let error = restart_error("0123456789abcdef0123456789abcdef", &executor).await;
706
707        assert!(
708            WorkerRestartLeftNoWorker::marks(&error),
709            "the worker was stopped and nothing replaced it: {error:#}"
710        );
711        assert!(
712            executor.executed.lock().expect("executed commands").len() > 1,
713            "the restart should have failed after its stop, not during it"
714        );
715    }
716
717    #[cfg(unix)]
718    #[tokio::test]
719    async fn checkpoint_restart_waits_for_upgrade_while_recovery_defers_and_another_session_runs() {
720        use std::sync::atomic::{AtomicUsize, Ordering};
721        struct FailingCommands(AtomicUsize);
722        impl CommandExecutor for FailingCommands {
723            fn execute(&self, _: &CommandSpec) -> Result<CommandOutput> {
724                self.0.fetch_add(1, Ordering::SeqCst);
725                Ok(CommandOutput {
726                    status: 1,
727                    stdout: Vec::new(),
728                    stderr: b"stop refused".to_vec(),
729                })
730            }
731        }
732        let id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
733        let upgrading = crate::worker_lifecycle::WorkerPermit::try_acquire(id, "worker upgrade")
734            .unwrap()
735            .unwrap();
736        let executor = FailingCommands(AtomicUsize::new(0));
737        let checkpoint = restart_error(id, &executor);
738        tokio::pin!(checkpoint);
739        assert!(
740            tokio::time::timeout(Duration::from_millis(50), &mut checkpoint)
741                .await
742                .is_err()
743        );
744        let recovery = crate::session_manager::WorkerRecoveryPlan {
745            source_target: mj_core::state::TargetLocator::LocalBare {
746                worker_root: "/unused".into(),
747            },
748            target: None,
749            workspace: None,
750            exit_record: None,
751            liveness_probe: CommandSpec::new("probe", std::iter::empty::<&str>()),
752            binary_refresh: None,
753            launch_refresh: None,
754            restart: targets::CommandPlan {
755                description: "restart".into(),
756                commands: vec![],
757            },
758        };
759        assert_eq!(
760            crate::session_manager::recover_worker_controlled(recovery, true, Some(id), &executor)
761                .unwrap(),
762            crate::session_manager::WorkerRecoveryOutcome::Suppressed
763        );
764        assert_eq!(executor.0.load(Ordering::SeqCst), 0);
765        let other = FailingCommands(AtomicUsize::new(0));
766        tokio::time::timeout(
767            Duration::from_secs(2),
768            restart_error("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", &other),
769        )
770        .await
771        .unwrap();
772        assert_eq!(other.0.load(Ordering::SeqCst), 1);
773        drop(upgrading);
774        let error = tokio::time::timeout(Duration::from_secs(2), checkpoint)
775            .await
776            .unwrap();
777        assert!(!WorkerRestartLeftNoWorker::marks(&error));
778        assert_eq!(executor.0.load(Ordering::SeqCst), 1);
779    }
780}