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 stop_worker_after_target_recovery(
240 executor,
241 restart.backend,
242 session_id,
243 restart.worker_root,
244 )
245 .context(restart.messages.stop)?;
246 self.start_installed_worker(session_id, executor, restart)
247 .await
248 }
249
250 pub(super) async fn start_installed_worker(
256 &self,
257 session_id: &str,
258 executor: &(impl CommandExecutor + Sync),
259 restart: InstalledWorkerRestart<'_>,
260 ) -> Result<StandaloneSession> {
261 let InstalledWorkerRestart {
262 backend,
263 worker_root,
264 reconnect,
265 launch,
266 prepared,
267 messages,
268 } = restart;
269 let mut connection = async {
272 if !prepared {
276 let binary = worker_binary_for(backend, executor)?;
277 replace_installed_worker_binary(executor, backend, session_id, &binary)
278 .context(messages.replace)?;
279 if let Some(launch) = launch {
280 replace_installed_worker_launch_config(executor, backend, session_id, launch)
281 .context("install the current Mjolnir worker launch configuration")?;
282 }
283 }
284 start_worker(executor, backend, worker_root).context(messages.start)?;
285 match connect_started_worker_with_timeout(
288 reconnect,
289 session_id,
290 executor,
291 backend,
292 worker_root,
293 WORKER_RESTART_TIMEOUT,
294 )
295 .await
296 {
297 Ok(connection) => Ok(connection),
298 Err(error) => Err(
299 worker_probe_diagnosis(executor, backend, worker_root, error)
300 .context(messages.connect),
301 ),
302 }
303 }
304 .await
305 .map_err(|error| error.context(WorkerRestartLeftNoWorker))?;
306 let project_memory = match self.project_memory_sync_target(session_id) {
307 Ok(target) => Some(target),
308 Err(error) => {
309 tracing::warn!(
310 session_id,
311 error = format!("{error:#}"),
312 "{}",
313 messages.project_memory
314 );
315 None
316 }
317 };
318 connection.set_project_memory_target(project_memory);
319 async {
322 let checkpoint_only = connection.sync().await?.operational.checkpoint_only;
323 if let Some(launch) = launch {
324 anyhow::ensure!(
325 checkpoint_only
326 == (launch.run_mode
327 == mj_core::worker_launch::WorkerRunMode::CheckpointOnly),
328 "restarted worker did not enter the requested execution mode"
329 );
330 }
331 if checkpoint_only {
332 return Ok(());
333 }
334 wait_for_native_session(&mut connection, executor)
335 .await
336 .context(messages.native_session)?;
337 wait_for_idle_projection(&mut connection, WORKER_RESTART_TIMEOUT)
338 .await
339 .context("wait for ACP to go idle after worker restart")
340 }
341 .await
342 .map_err(mark_if_transport_died)?;
343 Ok(connection)
344 }
345}
346
347async fn wait_for_idle_projection(relay: &mut StandaloneSession, timeout: Duration) -> Result<()> {
356 let deadline = tokio::time::Instant::now() + timeout;
357 let mut last_ordinal = None;
358 let mut stable_polls = 0_u8;
359 loop {
360 let snapshot = relay.sync().await?;
361 let ordinal = snapshot.operational.latest_ordinal;
362 let goal_active =
363 snapshot.operational.goal.synchronized() && snapshot.operational.goal.active();
364 let idle = snapshot.operational.native_session_is_ready()
365 && (!snapshot.operational.has_work_in_flight() || goal_active);
366 if idle && (goal_active || last_ordinal == Some(ordinal)) {
367 stable_polls = stable_polls.saturating_add(1);
368 if stable_polls >= 3 {
369 return Ok(());
370 }
371 } else {
372 stable_polls = 0;
373 }
374 last_ordinal = Some(ordinal);
375 if snapshot.operational.execution == RelayExecutionState::Closed {
376 bail!("ACP runtime stopped before becoming idle");
377 }
378 if tokio::time::Instant::now() >= deadline {
379 bail!(
380 "ACP runtime did not become idle after worker restart (execution={:?}, ordinal={ordinal})",
381 snapshot.operational.execution
382 );
383 }
384 tokio::time::sleep(Duration::from_millis(200)).await;
385 }
386}
387
388#[cfg(test)]
389mod tests {
390 #[test]
391 fn a_dead_transport_after_reconnect_marks_the_restart_as_leaving_no_worker() {
392 let died = anyhow::Error::new(crate::worker_client::RelayTransportDead::new(
393 "relay proxy disconnected during attach",
394 ))
395 .context("wait for ACP session after restarting the worker for checkpoint");
396 assert!(super::WorkerRestartLeftNoWorker::marks(
397 &super::mark_if_transport_died(died)
398 ));
399
400 let slow = anyhow::anyhow!("timed out waiting for the ACP session")
401 .context("wait for ACP session after restarting the worker for checkpoint");
402 let slow = super::mark_if_transport_died(slow);
403 assert!(!super::WorkerRestartLeftNoWorker::marks(&slow), "{slow:#}");
404 }
405
406 use super::*;
407
408 #[cfg(unix)]
409 use std::sync::Mutex;
410
411 #[cfg(unix)]
412 use crate::targets::CommandOutput;
413
414 #[cfg(unix)]
417 struct StopSucceedsThenFails {
418 executed: Mutex<Vec<String>>,
419 }
420
421 #[cfg(unix)]
422 impl CommandExecutor for StopSucceedsThenFails {
423 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
424 let mut executed = self.executed.lock().expect("executed commands");
425 executed.push(command.program.clone());
426 if executed.len() == 1 {
427 return Ok(CommandOutput {
428 status: 0,
429 stdout: Vec::new(),
430 stderr: Vec::new(),
431 });
432 }
433 Ok(CommandOutput {
434 status: 1,
435 stdout: Vec::new(),
436 stderr: b"no such target".to_vec(),
437 })
438 }
439 }
440
441 #[cfg(unix)]
442 struct FailingStop;
443
444 #[cfg(unix)]
445 impl CommandExecutor for FailingStop {
446 fn execute(&self, _command: &CommandSpec) -> Result<CommandOutput> {
447 Ok(CommandOutput {
448 status: 1,
449 stdout: Vec::new(),
450 stderr: b"permission denied".to_vec(),
451 })
452 }
453 }
454
455 #[cfg(unix)]
456 fn bare_restart_controller() -> Controller {
457 Controller {
458 config: mj_core::config::Config::default(),
459 state: mj_core::state::State::default(),
460 }
461 }
462
463 #[cfg(unix)]
464 async fn restart_error(
465 session_id: &str,
466 executor: &(impl CommandExecutor + Sync),
467 ) -> anyhow::Error {
468 let worker_root = format!("/tmp/mjolnir-restart-test/{session_id}");
469 let backend = targets::TargetLocator::LocalBare {
470 worker_root: worker_root.clone(),
471 };
472 let reconnect = CommandSpec::new("unused", std::iter::empty::<&str>());
473 let result = bare_restart_controller()
474 .restart_worker_with_installed_binary(
475 session_id,
476 executor,
477 InstalledWorkerRestart {
478 backend: &backend,
479 worker_root: &worker_root,
480 reconnect: &reconnect,
481 launch: None,
482 prepared: false,
483 messages: &RESTART_FOR_CHECKPOINT,
484 },
485 )
486 .await;
487 match result {
488 Ok(_) => panic!("a failing executor unexpectedly restarted the worker"),
489 Err(error) => error,
490 }
491 }
492
493 #[cfg(unix)]
494 #[tokio::test]
495 async fn a_restart_that_could_not_stop_the_worker_leaves_it_running() {
496 let error = restart_error("0123456789abcdef0123456789abcdef", &FailingStop).await;
497
498 assert!(
499 !WorkerRestartLeftNoWorker::marks(&error),
500 "a failed stop may leave the old worker alive: {error:#}"
501 );
502 }
503
504 #[cfg(unix)]
505 #[tokio::test]
506 async fn a_restart_that_stopped_the_worker_and_then_failed_is_marked() {
507 let executor = StopSucceedsThenFails {
508 executed: Mutex::new(Vec::new()),
509 };
510
511 let error = restart_error("0123456789abcdef0123456789abcdef", &executor).await;
512
513 assert!(
514 WorkerRestartLeftNoWorker::marks(&error),
515 "the worker was stopped and nothing replaced it: {error:#}"
516 );
517 assert!(
518 executor.executed.lock().expect("executed commands").len() > 1,
519 "the restart should have failed after its stop, not during it"
520 );
521 }
522
523 #[test]
526 fn only_a_matching_reported_build_counts_as_current() {
527 let installed = "a".repeat(64);
528
529 assert!(worker_runs_installed_build(Some(&installed), &installed));
530 assert!(!worker_runs_installed_build(
531 Some(&"b".repeat(64)),
532 &installed
533 ));
534 assert!(
535 !worker_runs_installed_build(None, &installed),
536 "a worker too old to report a build is older than this controller"
537 );
538 }
539}