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 if lease.connection_mut().protocol_version() >= 21 {
32 if !lease
33 .connection_mut()
34 .reserve_idle(command_id.clone())
35 .await?
36 {
37 lease.release();
38 return Ok(None);
39 }
40 } else {
41 lease
44 .connection_mut()
45 .submit(
46 command_id.clone(),
47 RelayCommand::BeginCheckpoint {
48 reason: Some("idle worker replacement".into()),
49 },
50 )
51 .await?;
52 }
53 loop {
54 let mut snapshot = lease.connection_mut().sync().await?;
55 let ready = checkpoint_barrier_is_ready(&snapshot, &command_id);
56 snapshot.operational.checkpoint_barrier = None;
57 if !snapshot.operational.safe_to_replace(harness) {
58 return Ok(None);
61 }
62 if ready {
63 return Ok(Some(Self {
64 lease,
65 command_id,
66 harness,
67 _worker_owner: worker_owner,
68 }));
69 }
70 tokio::time::sleep(Duration::from_millis(25)).await;
71 }
72 })
73 .await
74 .context("idle worker reservation timed out; worker was left running")?
75 }
76
77 pub(in crate::controller) async fn verify_for_upgrade(&mut self) -> Result<bool> {
78 let mut snapshot = self.lease.connection_mut().sync().await?;
79 if !checkpoint_barrier_is_ready(&snapshot, &self.command_id) {
80 return Ok(false);
81 }
82 snapshot.operational.checkpoint_barrier = None;
83 Ok(snapshot.operational.safe_to_replace(self.harness))
84 }
85
86 pub(in crate::controller) fn finish_replacement(mut self, connection: StandaloneSession) {
87 self.lease.replace_connection(connection);
88 self.lease.release();
89 }
90
91 pub async fn acquire(handle: &ManagedSessionHandle, harness: HarnessKind) -> Result<Self> {
92 let worker_owner = crate::worker_lifecycle::WorkerPermit::acquire(
93 handle.session_id(),
94 "workspace write",
95 &ProcessExecutor,
96 )
97 .await?;
98 tokio::time::timeout(Duration::from_secs(30), async {
99 let mut lease = handle.lease_connection().await?;
100 let snapshot = lease.connection_mut().sync().await?;
101 ensure!(
102 snapshot.operational.safe_to_replace(harness),
103 "session must be live and idle with no queued or background work"
104 );
105 let command_id = new_command_id("workspace-write")?;
106 lease
107 .connection_mut()
108 .submit(
109 command_id.clone(),
110 RelayCommand::BeginCheckpoint {
111 reason: Some("API workspace file write".into()),
112 },
113 )
114 .await?;
115 loop {
116 let snapshot = lease.connection_mut().sync().await?;
117 if checkpoint_barrier_is_ready(&snapshot, &command_id) {
118 let mut operation = Self {
119 lease,
120 command_id,
121 harness,
122 _worker_owner: worker_owner,
123 };
124 operation.verify().await?;
125 return Ok(operation);
126 }
127 ensure!(
131 !snapshot.operational.has_work_in_flight(),
132 "session started work before the file barrier was ready"
133 );
134 tokio::time::sleep(Duration::from_millis(25)).await;
135 }
136 })
137 .await
138 .context("session did not become available for a file write within 30 seconds")?
139 }
140
141 pub async fn verify(&mut self) -> Result<()> {
142 let mut snapshot = self.lease.connection_mut().sync().await?;
143 ensure!(
144 checkpoint_barrier_is_ready(&snapshot, &self.command_id),
145 "file write lost its workspace barrier"
146 );
147 snapshot.operational.checkpoint_barrier = None;
148 snapshot.operational.queued_prompts.clear();
151 ensure!(
152 snapshot.operational.safe_to_replace(self.harness),
153 "session is no longer idle for the file write"
154 );
155 Ok(())
156 }
157
158 pub async fn release(mut self) -> Result<()> {
159 tokio::time::timeout(Duration::from_secs(30), async {
160 self.lease
161 .connection_mut()
162 .submit(
163 new_command_id("workspace-release")?,
164 RelayCommand::ReleaseCheckpoint {
165 barrier_command_id: self.command_id.clone(),
166 },
167 )
168 .await?;
169 loop {
170 let snapshot = self.lease.connection_mut().sync().await?;
171 if snapshot.operational.checkpoint_barrier.as_deref() != Some(&self.command_id) {
172 return Ok::<_, anyhow::Error>(());
173 }
174 tokio::time::sleep(Duration::from_millis(25)).await;
175 }
176 })
177 .await
178 .context("release file write barrier timed out")??;
179 self.lease.release();
180 Ok(())
181 }
182}
183
184pub(super) fn checkpoint_barrier_is_ready(
185 snapshot: &ManagedSessionSnapshot,
186 command_id: &str,
187) -> bool {
188 snapshot.operational.checkpoint_barrier.as_deref() == Some(command_id)
189 && snapshot.operational.checkpoint_ready.is_some()
190}