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, start_worker_durably,
23 stop_worker_after_target_recovery, 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 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(¤t_launch)? != serde_json::to_value(&launch)? {
204 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 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 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 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 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 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 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 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 let mut connection = async {
419 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 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 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
509async 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 #[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 #[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}