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 prepare_managed_harness_for_upgrade(executor, &backend, session_id, &binary, &launch)
180 .context("prepare the current managed harness before replacing the worker")?;
181 stage_worker_binary_for_upgrade(executor, &backend, session_id, &binary)
182 .context("stage the current worker while the old worker remains available")?;
183 let handle = manager
184 .wait_for_session(session_id, UPGRADE_LEASE_TIMEOUT)
185 .await?;
186 let harness = self.state.sessions[session_id].harness_kind;
187 let Some(mut lease) =
188 super::IdleWorkspaceLease::acquire_for_upgrade(&handle, harness).await?
189 else {
190 return Ok(WorkerUpgradeOutcome::Deferred);
191 };
192 if !lease.verify_for_upgrade().await? {
193 return Ok(WorkerUpgradeOutcome::Deferred);
194 }
195 install_staged_worker_binary(executor, &backend, session_id)
196 .context("install the prepared worker under its idle reservation")?;
197 replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
198 .context("install the worker launch configuration under its idle reservation")?;
199
200 let restarted = self
201 .restart_worker_with_installed_binary(
202 session_id,
203 executor,
204 InstalledWorkerRestart {
205 backend: &backend,
206 worker_root: &worker_root,
207 reconnect: &reconnect,
208 launch: Some(&launch),
209 prepared: true,
210 messages: &RESTART_FOR_UPGRADE,
211 },
212 )
213 .await;
214 match restarted {
215 Ok(connection) => {
216 lease.finish_replacement(connection);
217 Ok(WorkerUpgradeOutcome::Upgraded { build: installed })
218 }
219 Err(error) => {
220 drop(lease);
223 Err(error)
224 }
225 }
226 }
227
228 pub(super) async fn restart_worker_with_installed_binary(
231 &self,
232 session_id: &str,
233 executor: &(impl CommandExecutor + Sync),
234 restart: InstalledWorkerRestart<'_>,
235 ) -> Result<StandaloneSession> {
236 let InstalledWorkerRestart {
237 backend,
238 worker_root,
239 reconnect,
240 launch,
241 prepared,
242 messages,
243 } = restart;
244 stop_worker_after_target_recovery(executor, backend, session_id, worker_root)
248 .context(messages.stop)?;
249 let mut connection = async {
252 if !prepared {
256 let binary = worker_binary_for(backend, executor)?;
257 replace_installed_worker_binary(executor, backend, session_id, &binary)
258 .context(messages.replace)?;
259 if let Some(launch) = launch {
260 replace_installed_worker_launch_config(executor, backend, session_id, launch)
261 .context("install the current Mjolnir worker launch configuration")?;
262 }
263 }
264 start_worker(executor, backend, worker_root).context(messages.start)?;
265 match connect_started_worker_with_timeout(
268 reconnect,
269 session_id,
270 executor,
271 backend,
272 worker_root,
273 WORKER_RESTART_TIMEOUT,
274 )
275 .await
276 {
277 Ok(connection) => Ok(connection),
278 Err(error) => Err(
279 worker_probe_diagnosis(executor, backend, worker_root, error)
280 .context(messages.connect),
281 ),
282 }
283 }
284 .await
285 .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
286 let project_memory = match self.project_memory_sync_target(session_id) {
287 Ok(target) => Some(target),
288 Err(error) => {
289 tracing::warn!(
290 session_id,
291 error = format!("{error:#}"),
292 "{}",
293 messages.project_memory
294 );
295 None
296 }
297 };
298 connection.set_project_memory_target(project_memory);
299 async {
302 let checkpoint_only = connection.sync().await?.operational.checkpoint_only;
303 if let Some(launch) = launch {
304 anyhow::ensure!(
305 checkpoint_only
306 == (launch.run_mode
307 == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
308 "restarted worker did not enter the requested execution mode"
309 );
310 }
311 if checkpoint_only {
312 return Ok(());
313 }
314 wait_for_native_session(&mut connection, executor)
315 .await
316 .context(messages.native_session)?;
317 wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT)
318 .await
319 .context("wait for ACP to go idle after worker restart")
320 }
321 .await
322 .map_err(mark_if_transport_died)?;
323 Ok(connection)
324 }
325}
326
327async fn wait_for_idle_projection(relay: &mut StandaloneSession, timeout: Duration) -> Result<()> {
336 let deadline = tokio::time::Instant::now() + timeout;
337 let mut last_ordinal = None;
338 let mut stable_polls = 0_u8;
339 loop {
340 let snapshot = relay.sync().await?;
341 let ordinal = snapshot.operational.latest_ordinal;
342 let goal_active =
343 snapshot.operational.goal.synchronized() && snapshot.operational.goal.active();
344 let idle = snapshot.operational.native_session_is_ready()
345 && (!snapshot.operational.has_work_in_flight() || goal_active);
346 if idle && (goal_active || last_ordinal == Some(ordinal)) {
347 stable_polls = stable_polls.saturating_add(1);
348 if stable_polls >= 3 {
349 return Ok(());
350 }
351 } else {
352 stable_polls = 0;
353 }
354 last_ordinal = Some(ordinal);
355 if snapshot.operational.execution == RelayExecutionState::Closed {
356 bail!("ACP runtime stopped before becoming idle");
357 }
358 if tokio::time::Instant::now() >= deadline {
359 bail!(
360 "ACP runtime did not become idle after worker restart (execution={:?}, ordinal={ordinal})",
361 snapshot.operational.execution
362 );
363 }
364 tokio::time::sleep(Duration::from_millis(200)).await;
365 }
366}
367
368#[cfg(test)]
369mod tests {
370 #[test]
371 fn a_dead_transport_after_reconnect_marks_the_restart_as_leaving_no_worker() {
372 let died = anyhow::Error::new(crate::worker_client::RelayTransportDead::new(
373 "relay proxy disconnected during attach",
374 ))
375 .context("wait for ACP session after restarting the worker for checkpoint");
376 assert!(super::WorkerRestartLeftNoWorker::marks(
377 &super::mark_if_transport_died(died)
378 ));
379
380 let slow = anyhow::anyhow!("timed out waiting for the ACP session")
381 .context("wait for ACP session after restarting the worker for checkpoint");
382 let slow = super::mark_if_transport_died(slow);
383 assert!(!super::WorkerRestartLeftNoWorker::marks(&slow), "{slow:#}");
384 }
385
386 use super::*;
387
388 #[cfg(unix)]
389 use std::sync::Mutex;
390
391 #[cfg(unix)]
392 use crate::targets::CommandOutput;
393
394 #[cfg(unix)]
397 struct StopSucceedsThenFails {
398 executed: Mutex<Vec<String>>,
399 }
400
401 #[cfg(unix)]
402 impl CommandExecutor for StopSucceedsThenFails {
403 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
404 let mut executed = self.executed.lock().expect("executed commands");
405 executed.push(command.program.clone());
406 if executed.len() == 1 {
407 return Ok(CommandOutput {
408 status: 0,
409 stdout: Vec::new(),
410 stderr: Vec::new(),
411 });
412 }
413 Ok(CommandOutput {
414 status: 1,
415 stdout: Vec::new(),
416 stderr: b"no such target".to_vec(),
417 })
418 }
419 }
420
421 #[cfg(unix)]
422 struct FailingStop;
423
424 #[cfg(unix)]
425 impl CommandExecutor for FailingStop {
426 fn execute(&self, _command: &CommandSpec) -> Result<CommandOutput> {
427 Ok(CommandOutput {
428 status: 1,
429 stdout: Vec::new(),
430 stderr: b"permission denied".to_vec(),
431 })
432 }
433 }
434
435 #[cfg(unix)]
436 fn bare_restart_controller() -> Controller {
437 Controller {
438 config: mj_core::config::Config::default(),
439 state: mj_core::state::State::default(),
440 }
441 }
442
443 #[cfg(unix)]
444 async fn restart_error(
445 session_id: &str,
446 executor: &(impl CommandExecutor + Sync),
447 ) -> anyhow::Error {
448 let worker_root = format!("/tmp/mjolnir-restart-test/{session_id}");
449 let backend = targets::TargetLocator::LocalBare {
450 worker_root: worker_root.clone(),
451 };
452 let reconnect = CommandSpec::new("unused", std::iter::empty::<&str>());
453 let result = bare_restart_controller()
454 .restart_worker_with_installed_binary(
455 session_id,
456 executor,
457 InstalledWorkerRestart {
458 backend: &backend,
459 worker_root: &worker_root,
460 reconnect: &reconnect,
461 launch: None,
462 prepared: false,
463 messages: &RESTART_FOR_CHECKPOINT,
464 },
465 )
466 .await;
467 match result {
468 Ok(_) => panic!("a failing executor unexpectedly restarted the worker"),
469 Err(error) => error,
470 }
471 }
472
473 #[cfg(unix)]
474 #[tokio::test]
475 async fn a_restart_that_could_not_stop_the_worker_leaves_it_running() {
476 let error = restart_error("0123456789abcdef0123456789abcdef", &FailingStop).await;
477
478 assert!(
479 !WorkerRestartLeftNoWorker::marks(&error),
480 "a failed stop may leave the old worker alive: {error:#}"
481 );
482 }
483
484 #[cfg(unix)]
485 #[tokio::test]
486 async fn a_restart_that_stopped_the_worker_and_then_failed_is_marked() {
487 let executor = StopSucceedsThenFails {
488 executed: Mutex::new(Vec::new()),
489 };
490
491 let error = restart_error("0123456789abcdef0123456789abcdef", &executor).await;
492
493 assert!(
494 WorkerRestartLeftNoWorker::marks(&error),
495 "the worker was stopped and nothing replaced it: {error:#}"
496 );
497 assert!(
498 executor.executed.lock().expect("executed commands").len() > 1,
499 "the restart should have failed after its stop, not during it"
500 );
501 }
502
503 #[test]
506 fn only_a_matching_reported_build_counts_as_current() {
507 let installed = "a".repeat(64);
508
509 assert!(worker_runs_installed_build(Some(&installed), &installed));
510 assert!(!worker_runs_installed_build(
511 Some(&"b".repeat(64)),
512 &installed
513 ));
514 assert!(
515 !worker_runs_installed_build(None, &installed),
516 "a worker too old to report a build is older than this controller"
517 );
518 }
519}