mj_controller/controller/checkpoint/
workspace_lease.rs1use super::*;
2
3pub struct IdleWorkspaceLease {
7 pub(super) lease: ManagedSessionLease,
8 pub(super) command_id: String,
9 pub(super) harness: HarnessKind,
10 _worker_owner: crate::worker_lifecycle::WorkerPermit,
11}
12
13impl IdleWorkspaceLease {
14 pub async fn acquire_for_upgrade(
16 handle: &ManagedSessionHandle,
17 harness: HarnessKind,
18 ) -> Result<Option<Self>> {
19 let Some(worker_owner) = crate::worker_lifecycle::WorkerPermit::try_acquire(
20 handle.session_id(),
21 "idle worker reservation",
22 )?
23 else {
24 return Ok(None);
25 };
26 tokio::time::timeout(Duration::from_secs(30), async {
27 let Some(mut lease) = handle.lease_idle_connection(harness).await? else {
28 return Ok(None);
29 };
30 let command_id = new_command_id("worker-upgrade")?;
31 let preparation_failed = lease
32 .connection_mut()
33 .operational()
34 .harness_preparation_failed();
35 if !preparation_failed && lease.connection_mut().protocol_version() >= 21 {
36 if !lease
37 .connection_mut()
38 .reserve_idle(command_id.clone())
39 .await?
40 {
41 lease.release();
42 return Ok(None);
43 }
44 } else {
45 lease
51 .connection_mut()
52 .submit(
53 command_id.clone(),
54 RelayCommand::BeginCheckpoint {
55 reason: Some("idle worker replacement".into()),
56 },
57 )
58 .await?;
59 }
60 loop {
61 let mut snapshot = lease.connection_mut().sync().await?;
62 let ready = checkpoint_barrier_is_ready(&snapshot, &command_id);
63 snapshot.operational.checkpoint_barrier = None;
64 if !snapshot.operational.safe_to_replace(harness) {
65 return Ok(None);
68 }
69 if ready {
70 return Ok(Some(Self {
71 lease,
72 command_id,
73 harness,
74 _worker_owner: worker_owner,
75 }));
76 }
77 tokio::time::sleep(Duration::from_millis(25)).await;
78 }
79 })
80 .await
81 .context("idle worker reservation timed out; worker was left running")?
82 }
83
84 pub(in crate::controller) async fn verify_for_upgrade(&mut self) -> Result<bool> {
85 let mut snapshot = self.lease.connection_mut().sync().await?;
86 if !checkpoint_barrier_is_ready(&snapshot, &self.command_id) {
87 return Ok(false);
88 }
89 snapshot.operational.checkpoint_barrier = None;
90 Ok(snapshot.operational.safe_to_replace(self.harness))
91 }
92
93 pub(in crate::controller) fn finish_replacement(mut self, connection: StandaloneSession) {
94 self.lease.replace_connection(connection);
95 self.lease.release();
96 }
97
98 pub async fn acquire(handle: &ManagedSessionHandle, harness: HarnessKind) -> Result<Self> {
99 let worker_owner = crate::worker_lifecycle::WorkerPermit::acquire(
100 handle.session_id(),
101 "workspace write",
102 &ProcessExecutor,
103 )
104 .await?;
105 tokio::time::timeout(Duration::from_secs(30), async {
106 let mut lease = handle.lease_connection().await?;
107 let snapshot = lease.connection_mut().sync().await?;
108 ensure!(
109 snapshot.operational.safe_to_replace(harness),
110 "session must be live and idle with no queued or background work"
111 );
112 let command_id = new_command_id("workspace-write")?;
113 lease
114 .connection_mut()
115 .submit(
116 command_id.clone(),
117 RelayCommand::BeginCheckpoint {
118 reason: Some("API workspace file write".into()),
119 },
120 )
121 .await?;
122 loop {
123 let snapshot = lease.connection_mut().sync().await?;
124 if checkpoint_barrier_is_ready(&snapshot, &command_id) {
125 let mut operation = Self {
126 lease,
127 command_id,
128 harness,
129 _worker_owner: worker_owner,
130 };
131 operation.verify().await?;
132 return Ok(operation);
133 }
134 ensure!(
138 !snapshot.operational.has_work_in_flight(),
139 "session started work before the file barrier was ready"
140 );
141 tokio::time::sleep(Duration::from_millis(25)).await;
142 }
143 })
144 .await
145 .context("session did not become available for a file write within 30 seconds")?
146 }
147
148 pub async fn verify(&mut self) -> Result<()> {
149 let mut snapshot = self.lease.connection_mut().sync().await?;
150 ensure!(
151 checkpoint_barrier_is_ready(&snapshot, &self.command_id),
152 "file write lost its workspace barrier"
153 );
154 snapshot.operational.checkpoint_barrier = None;
155 snapshot.operational.queued_prompts.clear();
158 ensure!(
159 snapshot.operational.safe_to_replace(self.harness),
160 "session is no longer idle for the file write"
161 );
162 Ok(())
163 }
164
165 pub async fn release(mut self) -> Result<()> {
166 tokio::time::timeout(Duration::from_secs(30), async {
167 self.lease
168 .connection_mut()
169 .submit(
170 new_command_id("workspace-release")?,
171 RelayCommand::ReleaseCheckpoint {
172 barrier_command_id: self.command_id.clone(),
173 },
174 )
175 .await?;
176 loop {
177 let snapshot = self.lease.connection_mut().sync().await?;
178 if snapshot.operational.checkpoint_barrier.as_deref() != Some(&self.command_id) {
179 return Ok::<_, anyhow::Error>(());
180 }
181 tokio::time::sleep(Duration::from_millis(25)).await;
182 }
183 })
184 .await
185 .context("release file write barrier timed out")??;
186 self.lease.release();
187 Ok(())
188 }
189}
190
191pub(super) fn checkpoint_barrier_is_ready(
192 snapshot: &ManagedSessionSnapshot,
193 command_id: &str,
194) -> bool {
195 snapshot.operational.checkpoint_barrier.as_deref() == Some(command_id)
196 && snapshot.operational.checkpoint_ready.is_some()
197}