1use 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
26const WORKER_RESTART_TIMEOUT: Duration = Duration::from_secs(300);
30
31const UPGRADE_LEASE_TIMEOUT: Duration = Duration::from_secs(5);
33
34#[derive(Debug)]
38pub struct WorkerRestartLeftNoWorker;
39
40impl WorkerRestartLeftNoWorker {
41 #[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
60fn 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
70pub(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
91pub(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
101const 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#[derive(Debug, Clone, PartialEq, Eq)]
115pub enum WorkerUpgradeOutcome {
116 Upgraded { build: String },
119 AlreadyCurrent { build: String },
121 Deferred,
124}
125
126impl WorkerUpgradeOutcome {
127 #[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
138fn worker_runs_installed_build(reported: Option<&str>, installed: &str) -> bool {
144 reported.is_some_and(|reported| reported == installed)
145}
146
147impl Controller {
148 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 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 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 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 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 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 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 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 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 let mut connection = async {
334 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 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 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
409async 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 #[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 #[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}