mj_controller/controller/
checkpoint.rs1use 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#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct SessionExportLayout {
55 pub backend: targets::TargetLocator,
57 pub workspace_root: String,
59 pub primary_repository: String,
62 pub repositories: Vec<CheckpointRepositorySpec>,
63 pub managed_worktree: Option<mj_core::state::ManagedWorktree>,
67}
68
69const CHECKPOINT_BARRIER_TIMEOUT: Duration = Duration::from_secs(30);
74const CHECKPOINT_CANCEL_TIMEOUT: Duration = Duration::from_secs(30);
78const CHECKPOINT_BARRIER_TIMEOUT_AFTER_RESTART: Duration = Duration::from_secs(300);
81
82pub(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;