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            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                // Failed preparation is terminal: no harness can start work,
46                // and old workers cannot include that fact in ReserveIdle.
47                // Their lifecycle dispatcher can still establish this barrier,
48                // which prevents independent shells from starting. Recheck live
49                // processes after it is ready. Historical workers use it too.
50                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                    // Dropping the connection releases the barrier even if its
66                    // acknowledgement was lost. No worker is stopped.
67                    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                // The bare flag misses a turn or a tool the projection has
135                // not caught up with. The lease already holds the barrier, so
136                // ask whether anything *else* is running.
137                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        // Commands may queue behind this barrier, but cannot begin until it
156        // releases. Their arrival does not invalidate an in-progress write.
157        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}