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::config::HarnessKind;
16use mj_core::relay::RelayExecutionState;
17
18use super::Controller;
19use super::readiness::{connect_started_worker_with_timeout, wait_for_native_session};
20use super::worker_binary::{
21 install_staged_worker_binary, prepare_managed_harness_for_upgrade,
22 replace_installed_worker_binary, replace_installed_worker_launch_config,
23 stage_worker_binary_for_upgrade, start_worker, start_worker_durably,
24 stop_worker_after_target_recovery, worker_binary_for, worker_probe_diagnosis,
25};
26
27const WORKER_RESTART_TIMEOUT: Duration = Duration::from_secs(300);
31
32const UPGRADE_LEASE_TIMEOUT: Duration = Duration::from_secs(5);
34
35#[derive(Debug)]
39pub struct WorkerRestartLeftNoWorker;
40
41impl WorkerRestartLeftNoWorker {
42 #[must_use]
48 pub fn marks(error: &anyhow::Error) -> bool {
49 error.downcast_ref::<Self>().is_some()
50 }
51}
52
53impl std::fmt::Display for WorkerRestartLeftNoWorker {
54 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
55 formatter.write_str("the worker restart left the session without a live worker")
56 }
57}
58
59impl std::error::Error for WorkerRestartLeftNoWorker {}
60
61fn mark_if_transport_died(error: anyhow::Error) -> anyhow::Error {
64 if crate::worker_client::RelayTransportDead::marks(&error) {
65 error.context(WorkerRestartLeftNoWorker)
66 } else {
67 error
68 }
69}
70
71pub(super) struct WorkerRestartMessages {
75 pub stop: &'static str,
76 pub replace: &'static str,
77 pub start: &'static str,
78 pub connect: &'static str,
79 pub project_memory: &'static str,
80 pub native_session: &'static str,
81}
82
83pub(super) struct InstalledWorkerRestart<'a> {
84 pub backend: &'a targets::TargetLocator,
85 pub worker_root: &'a str,
86 pub reconnect: &'a CommandSpec,
87 pub launch: Option<&'a mj_core::worker_launch::WorkerLaunchConfig>,
88 pub prepared: bool,
89 pub messages: &'a WorkerRestartMessages,
90}
91
92pub(super) const RESTART_FOR_CHECKPOINT: WorkerRestartMessages = WorkerRestartMessages {
94 stop: "stop wedged Mjolnir worker before retrying checkpoint",
95 replace: "replace Mjolnir worker binary before retrying checkpoint",
96 start: "start Mjolnir worker after interrupting a wedged ACP turn",
97 connect: "connect to Mjolnir worker after restarting it for checkpoint",
98 project_memory: "project memory will not be synchronized after checkpoint worker restart",
99 native_session: "wait for ACP session after restarting the worker for checkpoint",
100};
101
102const RESTART_FOR_UPGRADE: WorkerRestartMessages = WorkerRestartMessages {
105 stop: "stop the Mjolnir worker before installing the current binary",
106 replace: "install the current Mjolnir worker binary",
107 start: "start Mjolnir worker on the current binary",
108 connect: "connect to Mjolnir worker after upgrading its binary",
109 project_memory: "project memory will not be synchronized after the worker upgrade",
110 native_session: "wait for ACP session after upgrading the worker",
111};
112
113#[derive(Debug, Clone, PartialEq, Eq)]
116pub enum WorkerUpgradeOutcome {
117 Upgraded { build: String },
120 AlreadyCurrent { build: String },
122 Deferred,
125}
126
127impl WorkerUpgradeOutcome {
128 #[must_use]
131 pub fn build(&self) -> Option<&str> {
132 match self {
133 Self::Upgraded { build } | Self::AlreadyCurrent { build } => Some(build),
134 Self::Deferred => None,
135 }
136 }
137}
138
139fn worker_runs_installed_build(reported: Option<&str>, installed: &str) -> bool {
145 reported.is_some_and(|reported| reported == installed)
146}
147
148impl Controller {
149 pub async fn upgrade_session_worker(
158 &self,
159 session_id: &str,
160 executor: &(impl CommandExecutor + Sync),
161 manager: &SessionManagerControl,
162 reported_build: Option<&str>,
163 ) -> Result<WorkerUpgradeOutcome> {
164 let (backend, worker_root) = self.worker_placement(session_id)?;
165 let reconnect = targets::reconnect_plan(&backend, session_id)?
166 .commands
167 .into_iter()
168 .next()
169 .context("reconnect plan is empty")?;
170 let binary = worker_binary_for(&backend, executor)
171 .context("resolve the worker binary this controller would install")?;
172 let installed = mj_core::worker_launch::worker_executable_digest(&binary)?;
173 if worker_runs_installed_build(reported_build, &installed) {
174 return Ok(WorkerUpgradeOutcome::AlreadyCurrent { build: installed });
175 }
176
177 let launch = self.current_worker_launch_config(session_id, &backend)?;
180 if let Some(session) = self.state.sessions.get(session_id)
181 && session.build_cache.is_some()
182 {
183 self.prepare_build_cache_links(session, &backend, &launch, executor)
184 .context("prepare shared machine cache configuration before worker upgrade")?;
185 }
186 prepare_managed_harness_for_upgrade(executor, &backend, session_id, &binary, &launch)
187 .context("prepare the current managed harness before replacing the worker")?;
188 let staging = stage_worker_binary_for_upgrade(executor, &backend, session_id, &binary)
189 .context("stage the current worker while the old worker remains available")?;
190 let Some(owner) =
191 crate::worker_lifecycle::WorkerPermit::try_acquire(session_id, "worker upgrade")?
192 else {
193 return Ok(WorkerUpgradeOutcome::Deferred);
194 };
195 owner
196 .scope(async {
197 let target = self.state.sessions[session_id]
198 .target
199 .as_ref()
200 .context("worker upgrade has no target")?;
201 owner.verify_target(target)?;
202 let current = Controller::load()?;
203 let current_launch = current.current_worker_launch_config(session_id, &backend)?;
204 if serde_json::to_value(¤t_launch)? != serde_json::to_value(&launch)? {
205 return Ok(WorkerUpgradeOutcome::Deferred);
209 }
210 if crate::database::load_worker_restart(session_id)?.is_some() {
211 return Ok(WorkerUpgradeOutcome::Deferred);
212 }
213 let handle = manager
214 .wait_for_session(session_id, UPGRADE_LEASE_TIMEOUT)
215 .await?;
216 let harness = self.state.sessions[session_id].harness_kind;
217 let Ok(swap) = crate::upgrade::activity_unless_draining("worker swap") else {
220 return Ok(WorkerUpgradeOutcome::Deferred);
221 };
222 let Some(mut lease) =
223 super::IdleWorkspaceLease::acquire_for_upgrade(&handle, harness).await?
224 else {
225 return Ok(WorkerUpgradeOutcome::Deferred);
226 };
227 if !lease.verify_for_upgrade().await? {
228 return Ok(WorkerUpgradeOutcome::Deferred);
229 }
230 let operation_id = owner.operation_id().to_owned();
233 let intent = crate::database::WorkerRestartIntent {
234 operation_id: operation_id.clone(),
235 target: self.state.sessions[session_id]
236 .target
237 .clone()
238 .context("worker restart has no durable target")?,
239 desired_build: installed.clone(),
240 };
241 {
242 crate::database::begin_worker_restart(session_id, &intent)?;
243 install_staged_worker_binary(&owner, &staging, executor, &backend, session_id)
244 .context("install the prepared worker under its idle reservation")?;
245 replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
246 .context(
247 "install the worker launch configuration under its idle reservation",
248 )?;
249 crate::database::advance_worker_restart(
250 session_id,
251 &operation_id,
252 crate::database::WorkerRestartPhase::Swapping,
253 )?;
254 stop_worker_after_target_recovery(executor, &backend, session_id, &worker_root)
255 .context(RESTART_FOR_UPGRADE.stop)?;
256 start_worker_durably(&owner, target, executor, &backend, &worker_root)
257 .context(RESTART_FOR_UPGRADE.start)
258 .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
259 }
260 drop(swap);
263 let mut connection = connect_started_worker_with_timeout(
264 &reconnect,
265 session_id,
266 executor,
267 &backend,
268 &worker_root,
269 WORKER_RESTART_TIMEOUT,
270 )
271 .await
272 .context(RESTART_FOR_UPGRADE.connect)?;
273 anyhow::ensure!(
274 connection.snapshot().worker_build.as_deref() == Some(&installed),
275 "replacement worker reported an unexpected build"
276 );
277 anyhow::ensure!(
278 connection.snapshot().operational.checkpoint_only
279 == (launch.run_mode
280 == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
281 "replacement worker reported an unexpected execution mode"
282 );
283 let project_memory = match self.project_memory_sync_target(session_id) {
284 Ok(target) => Some(target),
285 Err(error) => {
286 tracing::warn!(
287 session_id,
288 error = format!("{error:#}"),
289 "project memory will not be synchronized after worker upgrade"
290 );
291 None
292 }
293 };
294 connection.set_project_memory_target(project_memory);
295 wait_for_native_session(&mut connection, executor, harness).await?;
296 wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT, executor).await?;
297 crate::database::finish_worker_restart(session_id, &operation_id)?;
298 lease.finish_replacement(connection);
299 Ok(WorkerUpgradeOutcome::Upgraded { build: installed })
300 })
301 .await
302 }
303
304 pub(super) async fn restart_worker_with_installed_binary(
307 &self,
308 session_id: &str,
309 executor: &(impl CommandExecutor + Sync),
310 restart: InstalledWorkerRestart<'_>,
311 ) -> Result<StandaloneSession> {
312 crate::worker_lifecycle::run(
313 session_id,
314 "restart worker with installed binary",
315 executor,
316 async {
317 let owner = crate::worker_lifecycle::require(session_id)?;
318 let target = self
319 .state
320 .sessions
321 .get(session_id)
322 .and_then(|session| session.target.as_ref());
323 if let Some(target) = target {
324 owner.verify_target(target)?;
325 if crate::database::load_worker_restart(session_id)?.is_some() {
326 let output =
327 executor.execute(&super::worker_binary::worker_liveness_command(
328 restart.backend,
329 restart.worker_root,
330 ))?;
331 anyhow::ensure!(
332 output.status == 0
333 && String::from_utf8_lossy(&output.stdout).trim() == "dead",
334 "worker replacement is still booting; reconnect before restarting it"
335 );
336 owner.settle_dead_restart(target)?;
337 }
338 owner.begin_restart(target, String::new())?;
339 crate::database::advance_worker_restart(
340 session_id,
341 owner.operation_id(),
342 crate::database::WorkerRestartPhase::Swapping,
343 )?;
344 }
345 stop_worker_after_target_recovery(
349 executor,
350 restart.backend,
351 session_id,
352 restart.worker_root,
353 )
354 .context(restart.messages.stop)?;
355 self.start_installed_worker(session_id, executor, restart)
356 .await
357 },
358 )
359 .await
360 }
361
362 pub(super) async fn start_installed_worker(
368 &self,
369 session_id: &str,
370 executor: &(impl CommandExecutor + Sync),
371 restart: InstalledWorkerRestart<'_>,
372 ) -> Result<StandaloneSession> {
373 crate::worker_lifecycle::run(session_id, "start installed worker", executor, async {
374 let InstalledWorkerRestart {
375 backend,
376 worker_root,
377 reconnect,
378 launch,
379 prepared,
380 messages,
381 } = restart;
382 let owner = crate::worker_lifecycle::require(session_id)?;
383 let session = self.state.sessions.get(session_id);
384 let harness = session
385 .map(|session| session.harness_kind)
386 .or_else(|| launch.map(|launch| launch.harness))
387 .unwrap_or(HarnessKind::Codex);
388 let observed_updated_at = session.map(|session| session.updated_at.clone());
389 let target = self
390 .state
391 .sessions
392 .get(session_id)
393 .and_then(|session| session.target.as_ref());
394 if let Some(target) = target {
395 owner.verify_target(target)?;
396 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
397 if intent.operation_id != owner.operation_id()
400 || crate::database::worker_restart_phase(session_id)?
401 == Some(crate::database::WorkerRestartPhase::AwaitingReadiness)
402 {
403 let output = executor.execute(
404 &super::worker_binary::worker_liveness_command(backend, worker_root),
405 )?;
406 anyhow::ensure!(
407 output.status == 0
408 && String::from_utf8_lossy(&output.stdout).trim() == "dead",
409 "worker replacement is still booting"
410 );
411 owner.settle_dead_restart(target)?;
412 }
413 }
414 if crate::database::load_worker_restart(session_id)?.is_none() {
415 owner.begin_restart(target, String::new())?;
416 crate::database::advance_worker_restart(
417 session_id,
418 owner.operation_id(),
419 crate::database::WorkerRestartPhase::Swapping,
420 )?;
421 }
422 }
423 let mut connection = async {
426 if !prepared {
430 let binary = worker_binary_for(backend, executor)?;
431 replace_installed_worker_binary(executor, backend, session_id, &binary)
432 .context(messages.replace)?;
433 if let Some(launch) = launch {
434 replace_installed_worker_launch_config(
435 executor, backend, session_id, launch,
436 )
437 .context("install the current Mjolnir worker launch configuration")?;
438 }
439 }
440 let started = match target {
441 Some(target) => {
442 start_worker_durably(&owner, target, executor, backend, worker_root)
443 }
444 None => start_worker(&owner, executor, backend, worker_root),
445 };
446 started.context(messages.start)?;
447 match connect_started_worker_with_timeout(
450 reconnect,
451 session_id,
452 executor,
453 backend,
454 worker_root,
455 WORKER_RESTART_TIMEOUT,
456 )
457 .await
458 {
459 Ok(connection) => Ok(connection),
460 Err(error) => {
461 Err(
462 worker_probe_diagnosis(executor, backend, worker_root, error)
463 .context(messages.connect),
464 )
465 }
466 }
467 }
468 .await
469 .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
470 let project_memory = match self.project_memory_sync_target(session_id) {
471 Ok(target) => Some(target),
472 Err(error) => {
473 tracing::warn!(
474 session_id,
475 error = format!("{error:#}"),
476 "{}",
477 messages.project_memory
478 );
479 None
480 }
481 };
482 connection.set_project_memory_target(project_memory);
483 let readiness = async {
486 let checkpoint_only = connection.sync().await?.operational.checkpoint_only;
487 if let Some(launch) = launch {
488 anyhow::ensure!(
489 checkpoint_only
490 == (launch.run_mode
491 == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
492 "restarted worker did not enter the requested execution mode"
493 );
494 }
495 if checkpoint_only {
496 return Ok(());
497 }
498 wait_for_native_session(&mut connection, executor, harness)
499 .await
500 .context(messages.native_session)?;
501 wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT, executor)
502 .await
503 .context("wait for ACP to go idle after worker restart")
504 }
505 .await
506 .map_err(mark_if_transport_died);
507 if let Err(error) = readiness {
508 if let (Some(failure), Some(observed_updated_at)) = (
509 error.downcast_ref::<super::readiness::HarnessPreparationFailure>(),
510 observed_updated_at.as_deref(),
511 ) {
512 self.persist_harness_preparation_failure(
513 session_id,
514 &failure.to_string(),
515 observed_updated_at,
516 )
517 .await?;
518 }
519 return Err(error);
520 }
521 if target.is_some() {
522 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
523 }
524 Ok(connection)
525 })
526 .await
527 }
528}
529
530async fn wait_for_idle_projection(
539 relay: &mut StandaloneSession,
540 timeout: Duration,
541 executor: &impl CommandExecutor,
542) -> Result<()> {
543 let deadline = tokio::time::Instant::now() + timeout;
544 let mut last_ordinal = None;
545 let mut stable_polls = 0_u8;
546 loop {
547 anyhow::ensure!(
548 !executor.cancellation_requested(),
549 "worker readiness wait cancelled"
550 );
551 let snapshot = relay.sync().await?;
552 let ordinal = snapshot.operational.latest_ordinal;
553 let goal_active =
554 snapshot.operational.goal.synchronized() && snapshot.operational.goal.active();
555 let idle = snapshot.operational.native_session_is_ready()
556 && (!snapshot.operational.has_work_in_flight() || goal_active);
557 if idle && (goal_active || last_ordinal == Some(ordinal)) {
558 stable_polls = stable_polls.saturating_add(1);
559 if stable_polls >= 3 {
560 return Ok(());
561 }
562 } else {
563 stable_polls = 0;
564 }
565 last_ordinal = Some(ordinal);
566 if snapshot.operational.execution == RelayExecutionState::Closed {
567 bail!("ACP runtime stopped before becoming idle");
568 }
569 if tokio::time::Instant::now() >= deadline {
570 bail!(
571 "ACP runtime did not become idle after worker restart (execution={:?}, ordinal={ordinal})",
572 snapshot.operational.execution
573 );
574 }
575 tokio::time::sleep(Duration::from_millis(200)).await;
576 }
577}
578
579#[cfg(test)]
580mod tests {
581 #[test]
583 fn a_dead_transport_after_reconnect_marks_the_restart_as_leaving_no_worker() {
584 let died = anyhow::Error::new(crate::worker_client::RelayTransportDead::new(
585 "relay proxy disconnected during attach",
586 ))
587 .context("wait for ACP session after restarting the worker for checkpoint");
588 assert!(super::WorkerRestartLeftNoWorker::marks(
589 &super::mark_if_transport_died(died)
590 ));
591
592 let slow = anyhow::anyhow!("timed out waiting for the ACP session")
593 .context("wait for ACP session after restarting the worker for checkpoint");
594 let slow = super::mark_if_transport_died(slow);
595 assert!(!super::WorkerRestartLeftNoWorker::marks(&slow), "{slow:#}");
596 }
597
598 #[cfg(unix)]
599 use super::*;
600
601 #[cfg(unix)]
602 use std::sync::Mutex;
603
604 #[cfg(unix)]
605 use crate::targets::CommandOutput;
606
607 #[cfg(unix)]
610 struct StopSucceedsThenFails {
611 executed: Mutex<Vec<String>>,
612 }
613
614 #[cfg(unix)]
615 impl CommandExecutor for StopSucceedsThenFails {
616 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
617 let mut executed = self.executed.lock().expect("executed commands");
618 executed.push(command.program.clone());
619 if executed.len() == 1 {
620 return Ok(CommandOutput {
621 status: 0,
622 stdout: Vec::new(),
623 stderr: Vec::new(),
624 });
625 }
626 Ok(CommandOutput {
627 status: 1,
628 stdout: Vec::new(),
629 stderr: b"no such target".to_vec(),
630 })
631 }
632 }
633
634 #[cfg(unix)]
635 struct FailingStop;
636
637 #[cfg(unix)]
638 impl CommandExecutor for FailingStop {
639 fn execute(&self, _command: &CommandSpec) -> Result<CommandOutput> {
640 Ok(CommandOutput {
641 status: 1,
642 stdout: Vec::new(),
643 stderr: b"permission denied".to_vec(),
644 })
645 }
646 }
647
648 #[cfg(unix)]
649 fn bare_restart_controller() -> Controller {
650 Controller {
651 config: mj_core::config::Config::default(),
652 state: mj_core::state::State::default(),
653 }
654 }
655
656 #[cfg(unix)]
657 async fn restart_error(
658 session_id: &str,
659 executor: &(impl CommandExecutor + Sync),
660 ) -> anyhow::Error {
661 let worker_root = format!("/tmp/mjolnir-restart-test/{session_id}");
662 let backend = targets::TargetLocator::LocalBare {
663 worker_root: worker_root.clone(),
664 };
665 let reconnect = CommandSpec::new("unused", std::iter::empty::<&str>());
666 let result = bare_restart_controller()
667 .restart_worker_with_installed_binary(
668 session_id,
669 executor,
670 InstalledWorkerRestart {
671 backend: &backend,
672 worker_root: &worker_root,
673 reconnect: &reconnect,
674 launch: None,
675 prepared: false,
676 messages: &RESTART_FOR_CHECKPOINT,
677 },
678 )
679 .await;
680 match result {
681 Ok(_) => panic!("a failing executor unexpectedly restarted the worker"),
682 Err(error) => error,
683 }
684 }
685
686 #[cfg(unix)]
687 #[tokio::test]
688 async fn a_restart_that_could_not_stop_the_worker_leaves_it_running() {
689 let error = restart_error("0123456789abcdef0123456789abcdef", &FailingStop).await;
690
691 assert!(
692 !WorkerRestartLeftNoWorker::marks(&error),
693 "a failed stop may leave the old worker alive: {error:#}"
694 );
695 }
696
697 #[cfg(unix)]
698 #[tokio::test]
700 async fn a_restart_that_stopped_the_worker_and_then_failed_is_marked() {
701 let executor = StopSucceedsThenFails {
702 executed: Mutex::new(Vec::new()),
703 };
704
705 let error = restart_error("0123456789abcdef0123456789abcdef", &executor).await;
706
707 assert!(
708 WorkerRestartLeftNoWorker::marks(&error),
709 "the worker was stopped and nothing replaced it: {error:#}"
710 );
711 assert!(
712 executor.executed.lock().expect("executed commands").len() > 1,
713 "the restart should have failed after its stop, not during it"
714 );
715 }
716
717 #[cfg(unix)]
718 #[tokio::test]
719 async fn checkpoint_restart_waits_for_upgrade_while_recovery_defers_and_another_session_runs() {
720 use std::sync::atomic::{AtomicUsize, Ordering};
721 struct FailingCommands(AtomicUsize);
722 impl CommandExecutor for FailingCommands {
723 fn execute(&self, _: &CommandSpec) -> Result<CommandOutput> {
724 self.0.fetch_add(1, Ordering::SeqCst);
725 Ok(CommandOutput {
726 status: 1,
727 stdout: Vec::new(),
728 stderr: b"stop refused".to_vec(),
729 })
730 }
731 }
732 let id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
733 let upgrading = crate::worker_lifecycle::WorkerPermit::try_acquire(id, "worker upgrade")
734 .unwrap()
735 .unwrap();
736 let executor = FailingCommands(AtomicUsize::new(0));
737 let checkpoint = restart_error(id, &executor);
738 tokio::pin!(checkpoint);
739 assert!(
740 tokio::time::timeout(Duration::from_millis(50), &mut checkpoint)
741 .await
742 .is_err()
743 );
744 let recovery = crate::session_manager::WorkerRecoveryPlan {
745 source_target: mj_core::state::TargetLocator::LocalBare {
746 worker_root: "/unused".into(),
747 },
748 target: None,
749 workspace: None,
750 exit_record: None,
751 liveness_probe: CommandSpec::new("probe", std::iter::empty::<&str>()),
752 binary_refresh: None,
753 launch_refresh: None,
754 restart: targets::CommandPlan {
755 description: "restart".into(),
756 commands: vec![],
757 },
758 };
759 assert_eq!(
760 crate::session_manager::recover_worker_controlled(recovery, true, Some(id), &executor)
761 .unwrap(),
762 crate::session_manager::WorkerRecoveryOutcome::Suppressed
763 );
764 assert_eq!(executor.0.load(Ordering::SeqCst), 0);
765 let other = FailingCommands(AtomicUsize::new(0));
766 tokio::time::timeout(
767 Duration::from_secs(2),
768 restart_error("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", &other),
769 )
770 .await
771 .unwrap();
772 assert_eq!(other.0.load(Ordering::SeqCst), 1);
773 drop(upgrading);
774 let error = tokio::time::timeout(Duration::from_secs(2), checkpoint)
775 .await
776 .unwrap();
777 assert!(!WorkerRestartLeftNoWorker::marks(&error));
778 assert_eq!(executor.0.load(Ordering::SeqCst), 1);
779 }
780}