mj_controller/controller/checkpoint/
barrier.rs1use super::*;
2
3pub(super) async fn connect_checkpoint_relay(
4 session_id: &str,
5 manager: Option<&SessionManagerControl>,
6 reconnect: &targets::CommandSpec,
7 project_memory: Option<crate::session_manager::ProjectMemorySyncTarget>,
8) -> Result<ControllerRelayLease> {
9 if let Some(manager) = manager {
10 let handle = manager
11 .wait_for_session(session_id, Duration::from_secs(5))
12 .await?;
13 let mut lease = handle.lease_connection().await?;
14 lease
15 .connection_mut()
16 .set_project_memory_target(project_memory);
17 Ok(ControllerRelayLease::Managed {
18 handle,
19 lease: Some(lease),
20 })
21 } else {
22 let target = crate::session_manager::RelaySessionTarget {
23 session_id: session_id.to_owned(),
24 spec: reconnect.clone(),
25 worker_recovery: None,
26 project_memory,
27 };
28 Ok(ControllerRelayLease::Standalone(
29 StandaloneSession::connect(&target).await?,
30 ))
31 }
32}
33
34pub(super) async fn adopt_restarted_checkpoint_relay(
35 session_id: &str,
36 manager: Option<&SessionManagerControl>,
37 connection: StandaloneSession,
38) -> Result<ControllerRelayLease> {
39 let Some(manager) = manager else {
40 return Ok(ControllerRelayLease::Standalone(connection));
41 };
42 let handle = manager
43 .wait_for_session(session_id, Duration::from_secs(5))
44 .await?;
45 match handle.lease_connection().await {
46 Ok(mut lease) => {
47 lease.replace_connection(connection);
48 Ok(ControllerRelayLease::Managed {
49 handle,
50 lease: Some(lease),
51 })
52 }
53 Err(error) => {
54 tracing::warn!(
55 session_id,
56 "session actor could not lease after worker restart; using the restarted proxy: {error:#}"
57 );
58 Ok(ControllerRelayLease::Standalone(connection))
59 }
60 }
61}
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub(super) enum BarrierBusyPolicy {
66 DeferWhileRunning,
72 InterruptWhileRunning,
76}
77
78impl BarrierBusyPolicy {
79 pub(super) fn of(exclusivity: LatchExclusivity) -> Self {
80 match exclusivity {
81 LatchExclusivity::ReleaseAfterLatch => Self::DeferWhileRunning,
82 LatchExclusivity::HoldThroughClose => Self::InterruptWhileRunning,
83 }
84 }
85}
86
87pub(super) async fn wait_for_checkpoint_barrier(
88 relay: &mut StandaloneSession,
89 session_id: &str,
90 command_id: &str,
91 timeout: Duration,
92 busy: BarrierBusyPolicy,
93 harness: HarnessKind,
94) -> Result<ManagedSessionSnapshot> {
95 let deadline = tokio::time::Instant::now() + timeout;
96 let mut cancel_submitted = false;
97 let mut cancel_deadline = None;
98 let mut cancel_started_at: Option<Instant> = None;
99 loop {
100 let snapshot = relay.sync().await?;
101 if busy == BarrierBusyPolicy::DeferWhileRunning
102 && let Some(wait) = snapshot.operational.routine_checkpoint_wait(harness)
103 {
104 return Err(CheckpointDeferred::from(wait).into());
110 }
111 if checkpoint_barrier_is_ready(&snapshot, command_id) {
112 if let Some(started_at) = cancel_started_at {
113 tracing::info!(
114 session_id,
115 barrier_command_id = command_id,
116 cancellation_ms = started_at.elapsed().as_millis() as u64,
117 "active turn cancellation settled before checkpoint barrier"
118 );
119 }
120 return Ok(snapshot);
121 }
122 if busy == BarrierBusyPolicy::InterruptWhileRunning
125 && snapshot.operational.activity_state().is_working()
126 && !cancel_submitted
127 {
128 let cancel_command_id = new_command_id("checkpoint-cancel-turn")?;
129 match relay
130 .submit(cancel_command_id, RelayCommand::CancelTurn)
131 .await
132 {
133 Ok(_) => {
134 cancel_submitted = true;
135 cancel_started_at = Some(Instant::now());
136 cancel_deadline = Some(tokio::time::Instant::now() + CHECKPOINT_CANCEL_TIMEOUT);
137 tracing::info!(
138 session_id,
139 barrier_command_id = command_id,
140 "requested active turn cancellation before checkpoint barrier"
141 );
142 }
143 Err(error) if checkpoint_cancel_turn_needs_worker_restart(&error) => {
144 return Err(error.context(
145 CheckpointBarrierUnreachable::cancel_turn_unavailable(
146 command_id,
147 relay.protocol_version(),
148 ),
149 ));
150 }
151 Err(error) if worker_connect_needs_restart(&error) => {
152 return Err(error.context(
153 CheckpointBarrierUnreachable::cancel_turn_unreachable(command_id),
154 ));
155 }
156 Err(error) => {
157 if let Ok(snapshot) = relay.sync().await
161 && checkpoint_barrier_is_ready(&snapshot, command_id)
162 {
163 tracing::info!(
164 session_id,
165 barrier_command_id = command_id,
166 "active turn settled while submitting checkpoint cancellation"
167 );
168 return Ok(snapshot);
169 }
170 return Err(error.context("cancel active ACP turn before checkpoint barrier"));
171 }
172 }
173 continue;
174 }
175 let out_of_time = tokio::time::Instant::now() >= cancel_deadline.unwrap_or(deadline);
176 if let Some(error) = checkpoint_barrier_wait_ended(
177 &snapshot,
178 command_id,
179 busy,
180 out_of_time,
181 cancel_submitted,
182 ) {
183 return Err(error);
184 }
185 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
186 }
187}
188
189pub(super) fn checkpoint_barrier_wait_ended(
199 snapshot: &ManagedSessionSnapshot,
200 command_id: &str,
201 busy: BarrierBusyPolicy,
202 out_of_time: bool,
203 cancel_submitted: bool,
204) -> Option<anyhow::Error> {
205 if snapshot.operational.execution == RelayExecutionState::Closed {
206 return Some(CheckpointBarrierUnreachable::runtime_stopped().into());
207 }
208 if busy == BarrierBusyPolicy::DeferWhileRunning {
209 if snapshot.operational.has_work_in_flight() {
210 return Some(CheckpointDeferred::harness_busy().into());
211 }
212 return out_of_time.then(|| CheckpointBarrierUnreachable::not_admitted(command_id).into());
213 }
214 if snapshot.operational.activity_state().is_working() {
217 return (out_of_time && cancel_submitted)
218 .then(|| CheckpointBarrierUnreachable::cancel_timed_out(command_id).into());
219 }
220 out_of_time.then(|| CheckpointBarrierUnreachable::not_admitted(command_id).into())
221}
222
223#[derive(Debug)]
230pub(super) struct CheckpointBarrierUnreachable(pub(super) String);
231
232impl CheckpointBarrierUnreachable {
233 pub(super) fn runtime_stopped() -> Self {
234 Self("ACP runtime stopped before reaching the checkpoint barrier".to_owned())
235 }
236
237 pub(super) fn not_admitted(command_id: &str) -> Self {
238 Self(format!(
239 "ACP relay did not reach checkpoint barrier {command_id}"
240 ))
241 }
242
243 pub(super) fn cancel_timed_out(command_id: &str) -> Self {
244 Self(format!(
245 "active ACP turn did not settle after cancellation before checkpoint barrier {command_id}"
246 ))
247 }
248
249 pub(super) fn cancel_turn_unavailable(command_id: &str, protocol_version: u32) -> Self {
250 Self(format!(
251 "worker protocol {protocol_version} cannot cancel the active ACP turn before checkpoint barrier {command_id} (requires protocol {})",
252 RelayCommand::CancelTurn.minimum_protocol(),
253 ))
254 }
255
256 pub(super) fn cancel_turn_unreachable(command_id: &str) -> Self {
257 Self(format!(
258 "worker transport became unavailable while cancelling the active ACP turn before checkpoint barrier {command_id}"
259 ))
260 }
261}
262
263impl std::fmt::Display for CheckpointBarrierUnreachable {
264 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
265 formatter.write_str(&self.0)
266 }
267}
268
269impl std::error::Error for CheckpointBarrierUnreachable {}
270
271pub(super) fn checkpoint_barrier_needs_worker_restart(error: &anyhow::Error) -> bool {
272 error
273 .downcast_ref::<CheckpointBarrierUnreachable>()
274 .is_some()
275}
276
277pub(super) fn checkpoint_cancel_turn_needs_worker_restart(error: &anyhow::Error) -> bool {
281 error.chain().any(|cause| {
282 let Some(rejected) = cause.downcast_ref::<RelayRejected>() else {
283 return false;
284 };
285 rejected.0.code == mj_core::relay::RelayErrorCode::IncompatibleProtocol
286 })
287}
288
289#[derive(Debug)]
299pub struct CheckpointDeferred(String);
300
301impl CheckpointDeferred {
302 pub fn harness_busy() -> Self {
303 Self("the agent is working; try again when it is idle".to_owned())
304 }
305
306 pub(super) fn background_work() -> Self {
307 Self("Kimi background-agent state could not be synchronized; checkpoint requires a synchronized empty task list".into())
308 }
309
310 pub(super) fn background_snapshot(
311 state: &mj_core::relay::RelayOperationalState,
312 harness: HarnessKind,
313 ) -> Self {
314 Self(
315 state
316 .checkpoint_background_blocker(harness)
317 .unwrap_or("background state changed during checkpoint")
318 .into(),
319 )
320 }
321
322 pub(super) fn frontier_moved() -> Self {
323 Self(
324 "the session moved past the checkpoint-ready cursor before the barrier latched, so this checkpoint was deferred"
325 .to_owned(),
326 )
327 }
328
329 pub(super) fn harness_turn_during_capture() -> Self {
330 Self(
331 "the agent started a turn of its own while target state was captured, so this checkpoint was deferred"
332 .to_owned(),
333 )
334 }
335}
336
337impl From<mj_core::activity::CheckpointWait> for CheckpointDeferred {
338 fn from(wait: mj_core::activity::CheckpointWait) -> Self {
339 match wait {
340 mj_core::activity::CheckpointWait::ProviderWork(reason) => Self(reason.to_owned()),
341 mj_core::activity::CheckpointWait::WorkInFlight => Self::harness_busy(),
342 }
343 }
344}
345
346impl std::fmt::Display for CheckpointDeferred {
347 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
348 formatter.write_str(&self.0)
349 }
350}
351
352impl std::error::Error for CheckpointDeferred {}
353
354pub fn checkpoint_was_deferred(error: &anyhow::Error) -> bool {
361 error.downcast_ref::<CheckpointDeferred>().is_some()
362}