Skip to main content

mj_controller/controller/checkpoint/
barrier.rs

1use 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/// What waiting for a barrier does while the session is working.
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub(super) enum BarrierBusyPolicy {
66    /// Give up as soon as the session is seen working. A checkpoint that can
67    /// run again later has nothing to gain from holding a barrier behind a
68    /// prompt or a turn the harness started on its own: the wait would only
69    /// end at the deadline, and the deadline means "wedged", which restarts
70    /// the worker and kills the work in flight.
71    DeferWhileRunning,
72    /// Request non-steering cancellation and wait for the turn to settle.
73    /// Close may interrupt work, but only an unresponsive or incompatible
74    /// worker needs restart recovery.
75    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            // Provider-owned work, a foreground tool, a turn the execution
105            // flag has not caught up with, or queued work can all appear
106            // after the initial sync and before the queued BeginCheckpoint is
107            // processed. Defer rather than let the deadline classify the
108            // worker as wedged and restart it underneath that work.
109            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        // A turn the execution flag has not caught up with still has to be
123        // cancelled, so this asks the shared state rather than the flag.
124        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                    // The turn can finish between the status sync and this
158                    // submit. If the barrier won that race, continue from its
159                    // durable ready state; otherwise preserve the rejection.
160                    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
189/// Why one sync of a barrier that is not ready yet ends the wait, or `None` to
190/// keep waiting.
191///
192/// The deadline means "wedged": it restarts the worker only after a close has
193/// already requested cancellation and the turn still has not settled. A
194/// checkpoint that can try again later defers as soon as it sees work. "Work"
195/// is the shared predicate, not the bare execution flag: a stale projection
196/// can report `Idle` while a harness turn or a foreground tool is still live,
197/// and treating that as wedged would restart the worker under it.
198pub(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    // Close already asked to interrupt the active turn, so only a turn that
215    // never settles after cancellation reaches the restart path.
216    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/// The ACP runtime never admitted a checkpoint barrier: it stopped first, or it
224/// never reached the barrier before the deadline.
225///
226/// [`wait_for_checkpoint_barrier`] is the only producer, and the retry decision
227/// downcasts for this marker rather than reading the message, so rewording a
228/// diagnostic cannot silently disable the restart-and-retry path.
229#[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
277/// A worker that cannot decode `CancelTurn` needs to be replaced before the
278/// close can retry the checkpoint with cancellation available. The relay client
279/// refuses the command for an older worker with the same code the worker uses.
280pub(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/// The session was working, so this checkpoint did not run. Nothing is wrong
290/// with the session, the target, or the last archive.
291///
292/// A busy session is the normal state of a session someone is using, including
293/// one working through a turn the harness started on its own after a
294/// background command. Treating that as a checkpoint failure would restart the
295/// worker, record a failure against the session, and back the next attempt off
296/// for hours. Callers that can try again later defer instead; the same work is
297/// copied at the next idle observation.
298#[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
354/// Whether a failed checkpoint only means the session was busy.
355///
356/// The marker is carried by the error, not by its text. It may be the root
357/// error or attached with `context`, and callers wrap checkpoint errors in
358/// further context. `anyhow`'s own downcast walks every context layer;
359/// `chain()` does not expose a context value, so it must not be used here.
360pub fn checkpoint_was_deferred(error: &anyhow::Error) -> bool {
361    error.downcast_ref::<CheckpointDeferred>().is_some()
362}