Skip to main content

mj_controller/controller/checkpoint/
workspace_lease.rs

1use super::*;
2
3/// An idle workspace operation holds the managed connection and a worker
4/// barrier. Dropping this value disconnects and cancels the barrier; releasing
5/// it resumes dispatch without claiming that an archive covers the journal.
6pub 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    /// Reserve an idle worker without taking a busy turn's control channel.
15    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                // Historical workers have the same disconnect-safe barrier.
42                // Never wait behind work or infer idle from a failed probe.
43                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                    // Dropping the connection releases the barrier even if its
59                    // acknowledgement was lost. No worker is stopped.
60                    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                // The bare flag misses a turn or a tool the projection has
128                // not caught up with. The lease already holds the barrier, so
129                // ask whether anything *else* is running.
130                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        // Commands may queue behind this barrier, but cannot begin until it
149        // releases. Their arrival does not invalidate an in-progress write.
150        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}