Skip to main content

mj_controller/controller/
checkpoint.rs

1//! Checkpoint export, latching, verification, and archive bookkeeping.
2
3use std::collections::BTreeSet;
4use std::ffi::OsStr;
5use std::path::{Path, PathBuf};
6use std::time::{Duration, Instant};
7
8use anyhow::{Context, Result, bail, ensure};
9
10use crate::checkpoint_transfer::{
11    CheckpointTransfer, capture_stdin_command, export_stdin_command, pack_stdin_command,
12};
13use crate::session_manager::{
14    ManagedSessionHandle, ManagedSessionLease, SessionManagerControl, StandaloneSession,
15    new_command_id, worker_connect_needs_restart,
16};
17use crate::worker_client::RelayRejected;
18use mj_checkpoint::archive::{
19    BundleManifest, CanonicalSessionSnapshot, SessionManifest, TargetManifest,
20    verify_archive_streaming,
21};
22use mj_checkpoint::checkpoint::{
23    CHECKPOINT_EXPORT_PROTOCOL_VERSION, CHECKPOINT_STAGING_PROTOCOL_VERSION, CapturedCheckpoint,
24    CheckpointCaptureSpec, CheckpointExportSpec, CheckpointPackSpec, CheckpointRepositoryCapture,
25    CheckpointRepositorySpec, NO_SESSION_ARTIFACTS, checkpoint_sha256,
26    current_native_session_received_prompt,
27};
28use mj_core::config::{HarnessKind, sessions_dir};
29use mj_core::state::{
30    CheckpointMetadata, ManagedSessionSnapshot, SessionRecord, SessionState, State,
31};
32use mj_transcript::projection::canonical_session_from_materialized;
33
34use crate::targets::{
35    self, CommandExecutor, CommandOutput, CommandSpec, ProcessExecutor, ProvisionStage,
36    ProvisionStageGuard,
37};
38use mj_core::relay::{RelayCommand, RelayCursor, RelayExecutionState};
39
40use super::backend::backend_locator;
41use super::readiness::wait_for_native_session_in_stage;
42use super::worker_restart::{InstalledWorkerRestart, RESTART_FOR_CHECKPOINT};
43use super::{
44    Controller, execute_checked, now, persist_session_record_transition_or_restore,
45    target_profile_home,
46};
47
48/// Where one session's work lives on its target.
49///
50/// Produced by [`Controller::session_export_layout`] and used both to build a
51/// checkpoint export specification and to run the worker's export commands
52/// against the right repository.
53#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct SessionExportLayout {
55    /// The provisioned target the session's commands run on.
56    pub backend: targets::TargetLocator,
57    /// The directory on the target that every repository is relative to.
58    pub workspace_root: String,
59    /// The id, within `repositories`, of the repository a caller means when it
60    /// names no repository.
61    pub primary_repository: String,
62    pub repositories: Vec<CheckpointRepositorySpec>,
63    /// Set when the session works in a Hel-owned worktree of the user's own
64    /// checkout, which is what records the branch and base commit an export
65    /// compares against.
66    pub managed_worktree: Option<mj_core::state::ManagedWorktree>,
67}
68
69/// How long an idle relay may fail to admit a barrier before its worker is
70/// treated as wedged. Busy recovery checkpoints defer immediately. A close
71/// sends a non-steering turn cancellation and gives the worker this same
72/// bounded interval to settle before recovery restarts it.
73const CHECKPOINT_BARRIER_TIMEOUT: Duration = Duration::from_secs(30);
74/// A close gets a fresh cancellation grace period once the worker accepts the
75/// request. This keeps an expensive status sync from consuming the whole
76/// cancellation budget before the worker has had a chance to settle.
77const CHECKPOINT_CANCEL_TIMEOUT: Duration = Duration::from_secs(30);
78/// After a wedged ACP forces a worker restart, wait as long as native-session
79/// startup: session/load of a long kimi transcript can outlast 30s.
80const CHECKPOINT_BARRIER_TIMEOUT_AFTER_RESTART: Duration = Duration::from_secs(300);
81
82/// Whether a checkpoint whose worker restart failed should start the worker
83/// once more without its harness. That restart only fails this way when the
84/// worker was stopped and none came back, often because the harness itself
85/// cannot start (R4-3). A checkpoint needs only the relay journal and the
86/// files on the target, so a suspend or move, which holds the latch through
87/// close and is ending the session anyway, can still save it. A routine
88/// recovery copy does not: the session would be left unable to run.
89pub(super) fn restart_falls_back_to_checkpoint_only(
90    exclusivity: LatchExclusivity,
91    error: &anyhow::Error,
92) -> bool {
93    exclusivity == LatchExclusivity::HoldThroughClose
94        && super::worker_restart::WorkerRestartLeftNoWorker::marks(error)
95}
96
97mod archives;
98pub use archives::*;
99mod lease;
100pub use lease::*;
101mod barrier;
102mod capture;
103mod latched;
104mod layout;
105mod persist;
106mod relay;
107pub use barrier::*;
108mod workspace_lease;
109pub use workspace_lease::*;
110mod validate;
111use validate::*;
112mod staging;
113pub(super) use staging::*;
114
115#[cfg(test)]
116pub(crate) mod tests;