Skip to main content

mj_controller/
worker_upgrade.rs

1//! Background worker-binary upgrade policy and coordination.
2//!
3//! A running session keeps the worker it started with, so a session that is
4//! never stopped never gains anything a newer daemon's worker learned. This
5//! coordinator watches the same session views the recovery coordinator does
6//! and, when a session is quiet and its worker is a different build from the
7//! one this controller would install, replaces it in place.
8//!
9//! Quiet is the whole safety argument: stopping a worker tears down the ACP
10//! bridge with it, so an upgrade may only run when nothing would be lost.
11
12use std::collections::BTreeMap;
13use std::sync::Arc;
14use std::sync::atomic::Ordering;
15use std::time::Duration;
16
17use chrono::{DateTime, Utc};
18use tokio::task::JoinSet;
19
20use crate::controller::{Controller, WorkerUpgradeOutcome};
21use crate::recovery::{backoff_delay, elapsed_at_least};
22use crate::recovery_gate::{
23    ObservationReceiver, ObservationSender, RecoveryGate, RecoveryObserver, observation_channel,
24};
25use crate::session_manager::SessionManagerControl;
26use crate::targets::CancellableProcessExecutor;
27use mj_core::config::Config;
28use mj_core::state::{SessionRecord, SessionState, State};
29
30/// How long a failed upgrade waits before it is tried again, doubling per
31/// consecutive failure. The daemon periodically reobserves quiet sessions, so
32/// without this one broken target would be retried continuously.
33const WORKER_UPGRADE_RETRY_INTERVAL: Duration = Duration::from_secs(10 * 60);
34const WORKER_UPGRADE_DEFERRED_INTERVAL: Duration = Duration::from_secs(30);
35
36/// Ceiling on the widening retry delay, so a target that is broken rather than
37/// blipping is still probed, just rarely.
38const MAX_WORKER_UPGRADE_RETRY_INTERVAL: Duration = Duration::from_secs(2 * 60 * 60);
39
40/// Upper bound on one upgrade. Stopping, installing, starting and waiting for
41/// a recovered journal all happen inside it; past this the attempt is a
42/// reported failure rather than a session whose upgrade never ends.
43const WORKER_UPGRADE_TIMEOUT: Duration = Duration::from_secs(15 * 60);
44
45/// One session view, reduced to what the upgrade decision reads.
46#[derive(Debug, Clone)]
47pub struct WorkerUpgradeObservation {
48    pub session: SessionRecord,
49    pub config: Config,
50    /// Content address the connected worker reported in hello, or `None` when
51    /// it is too old to report one. `None` counts as outdated.
52    pub worker_build: Option<String>,
53    /// Whether replacing the worker now would destroy nothing. See
54    /// [`mj_core::relay::RelayOperationalState::is_quiet`].
55    ///
56    /// An observer may skip reporting a session that is working - the daemon
57    /// does, to keep the config clone off the streaming path - but the rule
58    /// lives here, so a busy observation that does arrive is still refused.
59    pub quiet: bool,
60}
61
62/// Reports session activity to the upgrade coordinator.
63///
64/// Like the recovery observer, this is a queued hand-off: the caller is an
65/// event loop and must never wait on an upgrade decision.
66#[derive(Clone)]
67pub struct WorkerUpgradeObserver {
68    observations: ObservationSender<WorkerUpgradeObservation>,
69}
70
71impl WorkerUpgradeObserver {
72    pub fn observe(&self, observation: WorkerUpgradeObservation) {
73        self.observations
74            .send(observation.session.id.clone(), observation, |_, _| {});
75    }
76}
77
78#[derive(Debug, Clone)]
79pub struct WorkerUpgradeResult {
80    pub session_id: String,
81    pub outcome: Result<WorkerUpgradeOutcome, String>,
82    /// The attempt was preempted by a lifecycle operation or by coordinator
83    /// shutdown. It judged nothing, so it is neither a success nor a failure.
84    pub cancelled: bool,
85}
86
87pub struct WorkerUpgradeCoordinator {
88    supervisor: Option<tokio::task::JoinHandle<()>>,
89    observer: WorkerUpgradeObserver,
90    results: ObservationReceiver<WorkerUpgradeResult>,
91    gate: Arc<RecoveryGate>,
92}
93
94impl Drop for WorkerUpgradeCoordinator {
95    fn drop(&mut self) {
96        // Closing the shared gate prevents new attempts and cancels existing
97        // ones atomically. Both supervisors retain ownership until settlement.
98        self.gate.close();
99    }
100}
101
102impl WorkerUpgradeCoordinator {
103    /// Start the coordinator, sharing the recovery observer's gate so a
104    /// recovery copy and an upgrade never touch one session at the same time.
105    pub fn spawn(session_manager: SessionManagerControl, recovery: &RecoveryObserver) -> Self {
106        let (observations_tx, mut observations_rx) =
107            observation_channel::<WorkerUpgradeObservation>();
108        let (results_tx, results_rx) = observation_channel();
109        let gate = recovery.gate.clone();
110        let coordinator_gate = gate.clone();
111        let supervisor = tokio::spawn(async move {
112            let mut policies = BTreeMap::<String, PolicyState>::new();
113            let mut attempts = JoinSet::new();
114            let mut closing = false;
115            loop {
116                if closing && attempts.is_empty() {
117                    break;
118                }
119                tokio::select! {
120                    _ = coordinator_gate.closed(), if !closing => { closing = true; }
121                    observed = observations_rx.recv(), if !closing => {
122                        let Some(observation) = observed else { coordinator_gate.close(); closing = true; continue; };
123                        let session_id = observation.session.id.clone();
124                        let policy = policies.entry(session_id.clone()).or_default();
125                        policy.observe(&observation);
126                        if !policy.due(&observation, Utc::now()) {
127                            continue;
128                        }
129                        let Some(upgrade_cancelled) = coordinator_gate.try_start(&session_id)
130                        else {
131                            continue;
132                        };
133                        let session_manager = session_manager.clone();
134                        let task_cancelled = upgrade_cancelled.cancellation();
135                        let task_session_id = session_id.clone();
136                        attempts.spawn(async move {
137                            let (joined, admission) = upgrade_cancelled.run_blocking(move |cancelled| {
138                                let Some(session) = crate::recovery_gate::current_background_session(&observation.session)
139                                    .map_err(|error| format!("read current upgrade placement: {error:#}"))?
140                                else { return Ok(WorkerUpgradeOutcome::Deferred); };
141                                let mut state = State::default();
142                                state.sessions.insert(task_session_id.clone(), session);
143                                let controller = Controller {
144                                    config: observation.config,
145                                    state,
146                                };
147                                let executor = CancellableProcessExecutor::new(cancelled)
148                                    .with_deadline(WORKER_UPGRADE_TIMEOUT);
149                                mj_core::runtime::block_on(controller.upgrade_session_worker(
150                                    &task_session_id,
151                                    &executor,
152                                    &session_manager,
153                                    observation.worker_build.as_deref(),
154                                ))
155                                .and_then(|result| result)
156                                    .map_err(|error| format!("{error:#}"))
157                            })
158                            .await;
159                            let outcome = match joined {
160                                Ok(outcome) => outcome,
161                                Err(error) => Err(format!("worker upgrade task failed: {error}")),
162                            };
163                            let result = WorkerUpgradeResult {
164                                session_id,
165                                outcome,
166                                cancelled: task_cancelled.load(Ordering::Acquire),
167                            };
168                            (result, admission)
169                        });
170                    }
171                    completed = attempts.join_next(), if !attempts.is_empty() => {
172                        let Some(completed) = completed else { continue };
173                        let (result, _admission) = match completed {
174                            Ok(completed) => completed,
175                            Err(error) => { tracing::error!(%error, "worker upgrade attempt supervisor failed"); continue; }
176                        };
177                        let policy = policies.entry(result.session_id.clone()).or_default();
178                        policy.record(&result, Utc::now());
179                        if let Err(error) = &result.outcome && !result.cancelled {
180                            tracing::warn!(session_id = %result.session_id, %error, "background worker upgrade failed");
181                        }
182                        results_tx.send(result.session_id.clone(), result, |_, _| {});
183                    }
184                }
185            }
186        });
187        Self {
188            supervisor: Some(supervisor),
189            observer: WorkerUpgradeObserver {
190                observations: observations_tx,
191            },
192            results: results_rx,
193            gate,
194        }
195    }
196
197    /// Close admission and wait for every admitted executor and durable
198    /// settlement. The daemon applies its shared shutdown watchdog outside.
199    pub async fn shutdown(&mut self) -> anyhow::Result<()> {
200        self.gate.close();
201        if let Some(supervisor) = self.supervisor.take() {
202            supervisor.await.map_err(|error| {
203                anyhow::anyhow!("background coordinator supervisor failed: {error}")
204            })?;
205        }
206        Ok(())
207    }
208
209    pub fn observer(&self) -> WorkerUpgradeObserver {
210        self.observer.clone()
211    }
212
213    pub fn try_result(&mut self) -> Option<WorkerUpgradeResult> {
214        self.results.try_recv()
215    }
216}
217
218/// What the coordinator remembers about one session between observations.
219#[derive(Debug, Default, PartialEq, Eq)]
220struct PolicyState {
221    /// The build proved current for this coordinator. While the observed
222    /// build still matches it, no attempt is needed and nothing is hashed.
223    current_build: Option<String>,
224    failed_at: Option<DateTime<Utc>>,
225    deferred_at: Option<DateTime<Utc>>,
226    consecutive_failures: u32,
227}
228
229impl PolicyState {
230    /// Whether an upgrade attempt should start for this observation.
231    fn due(&self, observation: &WorkerUpgradeObservation, now: DateTime<Utc>) -> bool {
232        // Killing a worker mid-turn destroys the turn, and a session that is
233        // shutting down or already stopped has no worker worth replacing.
234        if !observation.quiet || observation.session.state != SessionState::Running {
235            return false;
236        }
237        if self.worker_is_known_current(observation.worker_build.as_deref()) {
238            return false;
239        }
240        if self.deferred_at.is_some_and(|deferred| {
241            !elapsed_at_least(deferred, now, WORKER_UPGRADE_DEFERRED_INTERVAL)
242        }) {
243            return false;
244        }
245        self.failed_at.is_none_or(|failed_at| {
246            elapsed_at_least(
247                failed_at,
248                now,
249                backoff_delay(
250                    WORKER_UPGRADE_RETRY_INTERVAL,
251                    MAX_WORKER_UPGRADE_RETRY_INTERVAL,
252                    self.consecutive_failures,
253                ),
254            )
255        })
256    }
257
258    /// Whether a previous attempt proved this exact build current. A worker
259    /// reporting no build is never current: it predates the field, so it
260    /// predates this controller.
261    fn worker_is_known_current(&self, worker_build: Option<&str>) -> bool {
262        let (Some(observed), Some(current)) = (worker_build, self.current_build.as_deref()) else {
263            return false;
264        };
265        observed == current
266    }
267
268    /// Fold in what one observation proves, before deciding whether to act.
269    ///
270    /// The one thing it can prove is that the last upgrade took: the worker
271    /// now reports the build that upgrade installed. That releases the
272    /// cooldown an upgrade leaves behind.
273    fn observe(&mut self, observation: &WorkerUpgradeObservation) {
274        if self.current_build.is_some()
275            && observation.worker_build.as_deref() == self.current_build.as_deref()
276        {
277            self.failed_at = None;
278            self.deferred_at = None;
279            self.consecutive_failures = 0;
280        }
281    }
282
283    fn record(&mut self, result: &WorkerUpgradeResult, now: DateTime<Utc>) {
284        if result.cancelled {
285            // A preempted attempt judged nothing: it must neither be counted
286            // as a failure nor suppress the next observation.
287            return;
288        }
289        self.deferred_at = None;
290        match &result.outcome {
291            Ok(WorkerUpgradeOutcome::Deferred) => {
292                // A busy worker or target recovery can refuse the swap. Do not
293                // repeat preparation on every unchanged quiet observation.
294                self.deferred_at = Some(now);
295            }
296            Ok(outcome @ WorkerUpgradeOutcome::AlreadyCurrent { .. }) => {
297                self.failed_at = None;
298                self.consecutive_failures = 0;
299                self.current_build = outcome.build().map(str::to_owned);
300            }
301            Ok(outcome @ WorkerUpgradeOutcome::Upgraded { .. }) => {
302                // The worker this session runs now, so the next observation
303                // reporting it needs no attempt - and, through `observe`,
304                // clears the cooldown below.
305                self.current_build = outcome.build().map(str::to_owned);
306                // An upgrade that did not take would otherwise restart this
307                // worker on every sync tick, forever. The same widening
308                // cooldown a failure gets bounds that; a worker that comes
309                // back reporting the installed build releases it immediately,
310                // so a healthy upgrade pays nothing.
311                self.consecutive_failures = self.consecutive_failures.saturating_add(1);
312                self.failed_at = Some(now);
313            }
314            Err(_) => {
315                self.consecutive_failures = self.consecutive_failures.saturating_add(1);
316                self.failed_at = Some(now);
317            }
318        }
319    }
320}
321
322#[cfg(test)]
323mod tests {
324    use super::*;
325
326    fn session_record(state: SessionState) -> SessionRecord {
327        SessionRecord {
328            project: None,
329            target_runtime: None,
330            launch_base: None,
331            launch_branch: None,
332            checkout: None,
333            publication: None,
334            build_cache: None,
335            container_workspace: None,
336            subagents: None,
337            create_managed_worktree: None,
338            workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
339            archived: false,
340            container_cpus: None,
341            container_memory: None,
342            id: "session-1".to_owned(),
343            title: "work".into(),
344            harness_kind: mj_core::config::HarnessKind::Codex,
345            last_profile: "codex-1".into(),
346            bundle_id: "hel".into(),
347            project_directory: None,
348            managed_worktree: None,
349            target_template_id: "podman".into(),
350            resource_allocation: None,
351            additional_mounts: Vec::new(),
352            state,
353            target: None,
354            native_session_id: None,
355            acp_session_title: None,
356            session_title_override: None,
357            created_at: "2026-08-09T12:00:00Z".into(),
358            updated_at: "2026-08-09T12:01:00Z".into(),
359            viewed_through_event_ordinal: 0,
360            draft_input: String::new(),
361            last_error: None,
362            last_checkpoint_error: None,
363            checkpoint: None,
364        }
365    }
366
367    fn observation(worker_build: Option<&str>, quiet: bool) -> WorkerUpgradeObservation {
368        WorkerUpgradeObservation {
369            session: session_record(SessionState::Running),
370            config: Config::default(),
371            worker_build: worker_build.map(str::to_owned),
372            quiet,
373        }
374    }
375
376    fn failure(detail: &str) -> WorkerUpgradeResult {
377        WorkerUpgradeResult {
378            session_id: "session-1".into(),
379            outcome: Err(detail.into()),
380            cancelled: false,
381        }
382    }
383
384    fn current(build: &str) -> WorkerUpgradeOutcome {
385        WorkerUpgradeOutcome::AlreadyCurrent {
386            build: build.to_owned(),
387        }
388    }
389
390    fn upgraded(build: &str) -> WorkerUpgradeOutcome {
391        WorkerUpgradeOutcome::Upgraded {
392            build: build.to_owned(),
393        }
394    }
395
396    fn success(outcome: WorkerUpgradeOutcome) -> WorkerUpgradeResult {
397        WorkerUpgradeResult {
398            session_id: "session-1".into(),
399            outcome: Ok(outcome),
400            cancelled: false,
401        }
402    }
403
404    /// The two facts that make an upgrade due, each on its own.
405    #[test]
406    fn only_a_quiet_session_with_an_unknown_build_is_due() {
407        let now = Utc::now();
408        let policy = PolicyState::default();
409
410        assert!(policy.due(&observation(Some("build-a"), true), now));
411        assert!(
412            !policy.due(&observation(Some("build-a"), false), now),
413            "a working session must not have its worker killed"
414        );
415        assert!(
416            policy.due(&observation(None, true), now),
417            "a worker too old to report a build is outdated"
418        );
419    }
420
421    #[test]
422    fn a_busy_turn_is_never_upgraded_no_matter_how_long_it_runs() {
423        let started = Utc::now();
424        let policy = PolicyState::default();
425        let two_days_later = started + chrono::Duration::days(2);
426
427        assert!(!policy.due(&observation(Some("old-build"), false), two_days_later));
428        assert!(
429            policy.due(&observation(Some("old-build"), true), two_days_later),
430            "the next quiet observation may upgrade without an age-based busy timeout"
431        );
432    }
433
434    /// A session that is closing, checkpointing or already stopped is not a
435    /// session whose worker may be replaced underneath it.
436    #[test]
437    fn only_a_running_session_is_due() {
438        let now = Utc::now();
439        let policy = PolicyState::default();
440        for state in [
441            SessionState::Provisioning,
442            SessionState::Disconnected,
443            SessionState::Checkpointing,
444            SessionState::Closing,
445            SessionState::Destroying,
446            SessionState::Stopped,
447            SessionState::Lost,
448            SessionState::Error,
449            SessionState::DestroyedWithDataLoss,
450        ] {
451            let mut observation = observation(Some("build-a"), true);
452            observation.session.state = state;
453            assert!(!policy.due(&observation, now), "{state:?}");
454        }
455    }
456
457    /// A proved build stays current until a different worker is observed.
458    #[test]
459    fn a_worker_proved_current_stays_trusted_for_coordinator_lifetime() {
460        let now = Utc::now();
461        let mut policy = PolicyState::default();
462        policy.record(&success(current("build-a")), now);
463
464        assert!(!policy.due(&observation(Some("build-a"), true), now));
465        assert!(
466            !policy.due(
467                &observation(Some("build-a"), true),
468                now + chrono::Duration::days(2)
469            ),
470            "the launched build remains trusted for the coordinator lifetime"
471        );
472        assert!(
473            policy.due(&observation(Some("build-b"), true), now),
474            "a different build is outdated however recently the last one was checked"
475        );
476    }
477
478    /// A failed upgrade waits, and waits longer each time, so a broken target
479    /// is not restarted on every sync tick.
480    #[test]
481    fn a_failed_upgrade_backs_off_and_widens() {
482        let now = Utc::now();
483        let mut policy = PolicyState::default();
484        policy.record(&failure("install the current Mjolnir worker binary"), now);
485
486        let interval = chrono::Duration::from_std(WORKER_UPGRADE_RETRY_INTERVAL).unwrap();
487        assert!(!policy.due(&observation(Some("build-a"), true), now));
488        assert!(!policy.due(
489            &observation(Some("build-a"), true),
490            now + interval - chrono::Duration::seconds(1)
491        ));
492        assert!(policy.due(&observation(Some("build-a"), true), now + interval));
493
494        policy.record(&failure("install the current Mjolnir worker binary"), now);
495        assert!(!policy.due(
496            &observation(Some("build-a"), true),
497            now + interval * 2 - chrono::Duration::seconds(1)
498        ));
499        assert!(policy.due(&observation(Some("build-a"), true), now + interval * 2));
500    }
501
502    /// A successful upgrade stops the session being probed again: the worker
503    /// that answers next is the new one, and confirming it clears the failure
504    /// run the upgrade itself was guarded by.
505    #[test]
506    fn a_successful_upgrade_stops_the_probing_and_confirming_it_clears_the_backoff() {
507        let now = Utc::now();
508        let mut policy = PolicyState::default();
509        policy.record(&failure("install the current Mjolnir worker binary"), now);
510        policy.record(&success(upgraded("build-b")), now);
511
512        let confirmed = observation(Some("build-b"), true);
513        policy.observe(&confirmed);
514        assert_eq!(policy.consecutive_failures, 0);
515        assert_eq!(policy.failed_at, None);
516        assert!(
517            !policy.due(&confirmed, now),
518            "the worker now runs the installed build, so nothing is due"
519        );
520    }
521
522    /// An upgrade that does not take - the worker comes back reporting a build
523    /// that is still not the installed one - must not restart that worker on
524    /// every freshness window forever.
525    #[test]
526    fn an_upgrade_that_does_not_take_backs_off_instead_of_looping() {
527        let now = Utc::now();
528        let mut policy = PolicyState::default();
529        policy.record(&success(upgraded("build-b")), now);
530
531        // The worker came back as something else, so nothing confirms the
532        // upgrade and the cooldown stands.
533        let unchanged = observation(Some("build-a"), true);
534        policy.observe(&unchanged);
535        let interval = chrono::Duration::from_std(WORKER_UPGRADE_RETRY_INTERVAL).unwrap();
536        assert!(!policy.due(&unchanged, now));
537        assert!(policy.due(&unchanged, now + interval));
538
539        policy.record(&success(upgraded("build-b")), now + interval);
540        policy.observe(&unchanged);
541        assert!(!policy.due(&unchanged, now + interval * 2));
542        assert!(policy.due(&unchanged, now + interval * 3));
543    }
544
545    /// Deferrals retain a bounded retry deadline even if observations stay quiet.
546    #[test]
547    fn a_deferred_upgrade_waits_before_repeating_preparation() {
548        let now = Utc::now();
549        let mut policy = PolicyState::default();
550        policy.record(&success(WorkerUpgradeOutcome::Deferred), now);
551
552        assert!(!policy.due(&observation(Some("build-a"), true), now));
553        let delay = chrono::Duration::from_std(WORKER_UPGRADE_DEFERRED_INTERVAL).unwrap();
554        assert!(!policy.due(
555            &observation(Some("build-a"), true),
556            now + delay - chrono::Duration::milliseconds(1)
557        ));
558        assert!(policy.due(&observation(Some("build-a"), true), now + delay));
559    }
560
561    /// A preempted attempt says nothing either way: it must not count as a
562    /// failure and must not delay the next attempt.
563    #[test]
564    fn a_preempted_attempt_is_neither_a_success_nor_a_failure() {
565        let now = Utc::now();
566        let mut policy = PolicyState::default();
567        policy.record(
568            &WorkerUpgradeResult {
569                session_id: "session-1".into(),
570                outcome: Err("operation cancelled".into()),
571                cancelled: true,
572            },
573            now,
574        );
575
576        assert_eq!(policy.consecutive_failures, 0);
577        assert!(policy.due(&observation(Some("build-a"), true), now));
578    }
579
580    /// Recovery holds the same per-session slot, so an upgrade cannot start
581    /// while a recovery copy is running for that session.
582    #[test]
583    fn the_shared_gate_keeps_an_upgrade_and_a_recovery_copy_apart() {
584        let gate = Arc::new(RecoveryGate::default());
585        let recovery_copy = gate.try_start("session-1").expect("the slot starts free");
586
587        assert!(gate.try_start("session-1").is_none());
588
589        drop(recovery_copy);
590        assert!(gate.try_start("session-1").is_some());
591    }
592
593    /// Observing must hand off and return: the daemon reports from its event
594    /// loop, and no upgrade decision may hold that loop up.
595    #[test]
596    fn observing_hands_off_without_waiting() {
597        let (observations, mut queued) = observation_channel();
598        let observer = WorkerUpgradeObserver { observations };
599
600        for _ in 0..64 {
601            observer.observe(observation(Some("build-a"), true));
602        }
603
604        let received = std::iter::from_fn(|| queued.try_recv()).count();
605        assert_eq!(received, 1);
606    }
607
608    /// A stopped coordinator leaves observing harmless.
609    #[test]
610    fn observing_a_stopped_coordinator_is_a_no_op() {
611        let (observations, queued) = observation_channel();
612        let observer = WorkerUpgradeObserver { observations };
613        drop(queued);
614
615        observer.observe(observation(Some("build-a"), true));
616    }
617}