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