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}
11
12impl IdleWorkspaceLease {
13    /// Reserve an idle worker without taking a busy turn's control channel.
14    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                // Historical workers have the same disconnect-safe barrier.
34                // Never wait behind work or infer idle from a failed probe.
35                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                    // Dropping the connection releases the barrier even if its
51                    // acknowledgement was lost. No worker is stopped.
52                    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                // The bare flag misses a turn or a tool the projection has
112                // not caught up with. The lease already holds the barrier, so
113                // ask whether anything *else* is running.
114                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        // Commands may queue behind this barrier, but cannot begin until it
133        // releases. Their arrival does not invalidate an in-progress write.
134        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}