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, stop_worker_after_target_recovery,
23    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        stage_worker_binary_for_upgrade(executor, &backend, session_id, &binary)
188            .context("stage the current worker while the old worker remains available")?;
189        let handle = manager
190            .wait_for_session(session_id, UPGRADE_LEASE_TIMEOUT)
191            .await?;
192        let harness = self.state.sessions[session_id].harness_kind;
193        // Reserve handoff admission before taking the worker's atomic idle
194        // reservation. A draining daemon must not take a worker connection.
195        let Ok(swap) = crate::upgrade::activity_unless_draining("worker swap") else {
196            return Ok(WorkerUpgradeOutcome::Deferred);
197        };
198        let Some(mut lease) =
199            super::IdleWorkspaceLease::acquire_for_upgrade(&handle, harness).await?
200        else {
201            return Ok(WorkerUpgradeOutcome::Deferred);
202        };
203        if !lease.verify_for_upgrade().await? {
204            return Ok(WorkerUpgradeOutcome::Deferred);
205        }
206        // The accepted swap records its target before touching a process. Once
207        // detached startup succeeds, the next daemon can resume observation.
208        let operation_id = crate::session_manager::new_command_id("worker-restart")?;
209        let intent = crate::database::WorkerRestartIntent {
210            operation_id: operation_id.clone(),
211            target: self.state.sessions[session_id]
212                .target
213                .clone()
214                .context("worker restart has no durable target")?,
215            desired_build: installed.clone(),
216        };
217        {
218            let target_lock = crate::recovery_gate::worker_target_mutex(session_id);
219            let _target = match target_lock.try_lock() {
220                Ok(target) => target,
221                Err(std::sync::TryLockError::WouldBlock) => {
222                    // Target recovery may be slow. Do not turn its work into
223                    // a daemon handoff blocker while waiting for ownership.
224                    return Ok(WorkerUpgradeOutcome::Deferred);
225                }
226                Err(std::sync::TryLockError::Poisoned(_)) => {
227                    bail!("worker target ownership lock poisoned");
228                }
229            };
230            crate::database::begin_worker_restart(session_id, &intent)?;
231            install_staged_worker_binary(executor, &backend, session_id)
232                .context("install the prepared worker under its idle reservation")?;
233            replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
234                .context("install the worker launch configuration under its idle reservation")?;
235            crate::database::advance_worker_restart(
236                session_id,
237                &operation_id,
238                crate::database::WorkerRestartPhase::Swapping,
239            )?;
240            stop_worker_after_target_recovery(executor, &backend, session_id, &worker_root)
241                .context(RESTART_FOR_UPGRADE.stop)?;
242            start_worker(executor, &backend, &worker_root)
243                .context(RESTART_FOR_UPGRADE.start)
244                .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
245            crate::database::advance_worker_restart(
246                session_id,
247                &operation_id,
248                crate::database::WorkerRestartPhase::AwaitingReadiness,
249            )?;
250        }
251        // The process now owns boot/journal recovery. Waiting for its socket
252        // is resumable, and must not hold daemon replacement for minutes.
253        drop(swap);
254        let mut connection = connect_started_worker_with_timeout(
255            &reconnect,
256            session_id,
257            executor,
258            &backend,
259            &worker_root,
260            WORKER_RESTART_TIMEOUT,
261        )
262        .await
263        .context(RESTART_FOR_UPGRADE.connect)?;
264        anyhow::ensure!(
265            connection.snapshot().worker_build.as_deref() == Some(&installed),
266            "replacement worker reported an unexpected build"
267        );
268        anyhow::ensure!(
269            connection.snapshot().operational.checkpoint_only
270                == (launch.run_mode == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
271            "replacement worker reported an unexpected execution mode"
272        );
273        let project_memory = match self.project_memory_sync_target(session_id) {
274            Ok(target) => Some(target),
275            Err(error) => {
276                tracing::warn!(
277                    session_id,
278                    error = format!("{error:#}"),
279                    "project memory will not be synchronized after worker upgrade"
280                );
281                None
282            }
283        };
284        connection.set_project_memory_target(project_memory);
285        crate::database::finish_worker_restart(session_id, &operation_id)?;
286        lease.finish_replacement(connection);
287        Ok(WorkerUpgradeOutcome::Upgraded { build: installed })
288    }
289
290    /// Stop the worker, install the binary this controller would provision,
291    /// start it and reconnect to the session it recovers.
292    pub(super) async fn restart_worker_with_installed_binary(
293        &self,
294        session_id: &str,
295        executor: &(impl CommandExecutor + Sync),
296        restart: InstalledWorkerRestart<'_>,
297    ) -> Result<StandaloneSession> {
298        // A failed stop may leave the old worker alive, so it stays outside the
299        // marker the start below applies: only steps after a successful stop
300        // can leave the session with no worker at all.
301        stop_worker_after_target_recovery(
302            executor,
303            restart.backend,
304            session_id,
305            restart.worker_root,
306        )
307        .context(restart.messages.stop)?;
308        self.start_installed_worker(session_id, executor, restart)
309            .await
310    }
311
312    /// The part of a restart after its stop: install the binary unless it is
313    /// `prepared`, start the worker on the existing worker root, connect with
314    /// the long restart timeout, and wait until its harness has loaded its
315    /// native session and gone idle. A parked sub-agent is started again with
316    /// exactly this sequence, since its worker was stopped when it was parked.
317    pub(super) async fn start_installed_worker(
318        &self,
319        session_id: &str,
320        executor: &(impl CommandExecutor + Sync),
321        restart: InstalledWorkerRestart<'_>,
322    ) -> Result<StandaloneSession> {
323        let InstalledWorkerRestart {
324            backend,
325            worker_root,
326            reconnect,
327            launch,
328            prepared,
329            messages,
330        } = restart;
331        // Everything up to the first successful connection either fails with no
332        // worker running or cannot tell: the marker covers all of it.
333        let mut connection = async {
334            // Copy through hel.next and rename. scp/cp onto a still-mapped hel
335            // fails with ETXTBSY ("dest open ... Failure") even after SIGKILL,
336            // and prepare_worker_files writes that path in place.
337            if !prepared {
338                let binary = worker_binary_for(backend, executor)?;
339                replace_installed_worker_binary(executor, backend, session_id, &binary)
340                    .context(messages.replace)?;
341                if let Some(launch) = launch {
342                    replace_installed_worker_launch_config(executor, backend, session_id, launch)
343                        .context("install the current Mjolnir worker launch configuration")?;
344                }
345            }
346            start_worker(executor, backend, worker_root).context(messages.start)?;
347            // Journal recovery runs before the daemon binds control.sock. A long
348            // kimi session can take well over the ordinary 30s startup window.
349            match connect_started_worker_with_timeout(
350                reconnect,
351                session_id,
352                executor,
353                backend,
354                worker_root,
355                WORKER_RESTART_TIMEOUT,
356            )
357            .await
358            {
359                Ok(connection) => Ok(connection),
360                Err(error) => Err(
361                    worker_probe_diagnosis(executor, backend, worker_root, error)
362                        .context(messages.connect),
363                ),
364            }
365        }
366        .await
367        .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
368        let project_memory = match self.project_memory_sync_target(session_id) {
369            Ok(target) => Some(target),
370            Err(error) => {
371                tracing::warn!(
372                    session_id,
373                    error = format!("{error:#}"),
374                    "{}",
375                    messages.project_memory
376                );
377                None
378            }
379        };
380        connection.set_project_memory_target(project_memory);
381        // A worker answered, so a failure from here on only means "no worker"
382        // when the transport to it died again.
383        async {
384            let checkpoint_only = connection.sync().await?.operational.checkpoint_only;
385            if let Some(launch) = launch {
386                anyhow::ensure!(
387                    checkpoint_only
388                        == (launch.run_mode
389                            == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
390                    "restarted worker did not enter the requested execution mode"
391                );
392            }
393            if checkpoint_only {
394                return Ok(());
395            }
396            wait_for_native_session(&mut connection, executor)
397                .await
398                .context(messages.native_session)?;
399            wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT)
400                .await
401                .context("wait for ACP to go idle after worker restart")
402        }
403        .await
404        .map_err(mark_if_transport_died)?;
405        Ok(connection)
406    }
407}
408
409/// Wait until a restarted worker's projection stops moving and reports idle.
410///
411/// Three stable polls, not one: a worker that has just recovered its journal
412/// can report idle between two events it is still applying. "Idle" is the
413/// shared predicate, not the bare execution flag, so a foreground tool or a
414/// turn the projection has not caught up with keeps the restart from being
415/// declared ready underneath it. A synchronized active goal is the one
416/// deliberate exception: the restarted worker is meant to continue it.
417async fn wait_for_idle_projection(relay: &mut StandaloneSession, timeout: Duration) -> Result<()> {
418    let deadline = tokio::time::Instant::now() + timeout;
419    let mut last_ordinal = None;
420    let mut stable_polls = 0_u8;
421    loop {
422        let snapshot = relay.sync().await?;
423        let ordinal = snapshot.operational.latest_ordinal;
424        let goal_active =
425            snapshot.operational.goal.synchronized() && snapshot.operational.goal.active();
426        let idle = snapshot.operational.native_session_is_ready()
427            && (!snapshot.operational.has_work_in_flight() || goal_active);
428        if idle && (goal_active || last_ordinal == Some(ordinal)) {
429            stable_polls = stable_polls.saturating_add(1);
430            if stable_polls >= 3 {
431                return Ok(());
432            }
433        } else {
434            stable_polls = 0;
435        }
436        last_ordinal = Some(ordinal);
437        if snapshot.operational.execution == RelayExecutionState::Closed {
438            bail!("ACP runtime stopped before becoming idle");
439        }
440        if tokio::time::Instant::now() >= deadline {
441            bail!(
442                "ACP runtime did not become idle after worker restart (execution={:?}, ordinal={ordinal})",
443                snapshot.operational.execution
444            );
445        }
446        tokio::time::sleep(Duration::from_millis(200)).await;
447    }
448}
449
450#[cfg(test)]
451mod tests {
452    #[test]
453    fn a_dead_transport_after_reconnect_marks_the_restart_as_leaving_no_worker() {
454        let died = anyhow::Error::new(crate::worker_client::RelayTransportDead::new(
455            "relay proxy disconnected during attach",
456        ))
457        .context("wait for ACP session after restarting the worker for checkpoint");
458        assert!(super::WorkerRestartLeftNoWorker::marks(
459            &super::mark_if_transport_died(died)
460        ));
461
462        let slow = anyhow::anyhow!("timed out waiting for the ACP session")
463            .context("wait for ACP session after restarting the worker for checkpoint");
464        let slow = super::mark_if_transport_died(slow);
465        assert!(!super::WorkerRestartLeftNoWorker::marks(&slow), "{slow:#}");
466    }
467
468    use super::*;
469
470    #[cfg(unix)]
471    use std::sync::Mutex;
472
473    #[cfg(unix)]
474    use crate::targets::CommandOutput;
475
476    /// Fails every command after the first, so a restart gets past its stop and
477    /// then loses the worker it was replacing.
478    #[cfg(unix)]
479    struct StopSucceedsThenFails {
480        executed: Mutex<Vec<String>>,
481    }
482
483    #[cfg(unix)]
484    impl CommandExecutor for StopSucceedsThenFails {
485        fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
486            let mut executed = self.executed.lock().expect("executed commands");
487            executed.push(command.program.clone());
488            if executed.len() == 1 {
489                return Ok(CommandOutput {
490                    status: 0,
491                    stdout: Vec::new(),
492                    stderr: Vec::new(),
493                });
494            }
495            Ok(CommandOutput {
496                status: 1,
497                stdout: Vec::new(),
498                stderr: b"no such target".to_vec(),
499            })
500        }
501    }
502
503    #[cfg(unix)]
504    struct FailingStop;
505
506    #[cfg(unix)]
507    impl CommandExecutor for FailingStop {
508        fn execute(&self, _command: &CommandSpec) -> Result<CommandOutput> {
509            Ok(CommandOutput {
510                status: 1,
511                stdout: Vec::new(),
512                stderr: b"permission denied".to_vec(),
513            })
514        }
515    }
516
517    #[cfg(unix)]
518    fn bare_restart_controller() -> Controller {
519        Controller {
520            config: mj_core::config::Config::default(),
521            state: mj_core::state::State::default(),
522        }
523    }
524
525    #[cfg(unix)]
526    async fn restart_error(
527        session_id: &str,
528        executor: &(impl CommandExecutor + Sync),
529    ) -> anyhow::Error {
530        let worker_root = format!("/tmp/mjolnir-restart-test/{session_id}");
531        let backend = targets::TargetLocator::LocalBare {
532            worker_root: worker_root.clone(),
533        };
534        let reconnect = CommandSpec::new("unused", std::iter::empty::<&str>());
535        let result = bare_restart_controller()
536            .restart_worker_with_installed_binary(
537                session_id,
538                executor,
539                InstalledWorkerRestart {
540                    backend: &backend,
541                    worker_root: &worker_root,
542                    reconnect: &reconnect,
543                    launch: None,
544                    prepared: false,
545                    messages: &RESTART_FOR_CHECKPOINT,
546                },
547            )
548            .await;
549        match result {
550            Ok(_) => panic!("a failing executor unexpectedly restarted the worker"),
551            Err(error) => error,
552        }
553    }
554
555    #[cfg(unix)]
556    #[tokio::test]
557    async fn a_restart_that_could_not_stop_the_worker_leaves_it_running() {
558        let error = restart_error("0123456789abcdef0123456789abcdef", &FailingStop).await;
559
560        assert!(
561            !WorkerRestartLeftNoWorker::marks(&error),
562            "a failed stop may leave the old worker alive: {error:#}"
563        );
564    }
565
566    #[cfg(unix)]
567    #[tokio::test]
568    async fn a_restart_that_stopped_the_worker_and_then_failed_is_marked() {
569        let executor = StopSucceedsThenFails {
570            executed: Mutex::new(Vec::new()),
571        };
572
573        let error = restart_error("0123456789abcdef0123456789abcdef", &executor).await;
574
575        assert!(
576            WorkerRestartLeftNoWorker::marks(&error),
577            "the worker was stopped and nothing replaced it: {error:#}"
578        );
579        assert!(
580            executor.executed.lock().expect("executed commands").len() > 1,
581            "the restart should have failed after its stop, not during it"
582        );
583    }
584
585    /// The three answers hello can produce, and what each means for the
586    /// worker's binary.
587    #[test]
588    fn only_a_matching_reported_build_counts_as_current() {
589        let installed = "a".repeat(64);
590
591        assert!(worker_runs_installed_build(Some(&installed), &installed));
592        assert!(!worker_runs_installed_build(
593            Some(&"b".repeat(64)),
594            &installed
595        ));
596        assert!(
597            !worker_runs_installed_build(None, &installed),
598            "a worker too old to report a build is older than this controller"
599        );
600    }
601}