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