Skip to main content

mj_controller/daemon/
resume.rs

1use super::*;
2
3impl RuntimeState {
4    /// Resume a session, and return nothing.
5    ///
6    /// This used to answer with the whole `MaterializedSession`. That reply
7    /// travels as one JSON frame against `MAX_FRAME_BYTES`, so a session whose
8    /// projection outgrew 8 MiB could not be resumed at all — it built a
9    /// several-hundred-megabyte buffer and then refused to send it. The
10    /// projection is already durable; a viewer reads it from the store.
11    pub async fn resume_session(self: &Arc<Self>, request: ResumeSessionRequest) -> Result<()> {
12        let _admission = self
13            .workspace_resume_gate(&request.workspace_id)
14            .read_owned()
15            .await;
16        let session_id = request.session_id.clone();
17        self.wait_for_deferred_cleanup(&session_id).await?;
18        // Whether it is already running is a boolean. Answering it used to
19        // load the entire projection so it could be handed back as the reply.
20        let already_running = blocking({
21            let session_id = session_id.clone();
22            move || {
23                let controller = Controller::load()?;
24                Ok(controller
25                    .state
26                    .sessions
27                    .get(&session_id)
28                    .is_some_and(|session| session.state == SessionState::Running))
29            }
30        })
31        .await?;
32        if already_running {
33            return Ok(());
34        }
35        let profile_id = request.profile_id.clone();
36        let target_template_id = request.target_template_id.clone();
37        let workspace_id = request.workspace_id.clone();
38        let rebind_workspace_id = workspace_id.clone();
39        let operation_session_id = session_id.clone();
40        let result = self.start_or_join_lifecycle_for_workspace(
41            session_id,
42            LifecycleKind::Resume,
43            Some(workspace_id),
44            move |state, session_id, cancelled| async move {
45                let _recovery_reservation = tokio::task::spawn_blocking({
46                    let observer = state.recovery_observer.clone();
47                    let session_id = session_id.clone();
48                    let cancelled = cancelled.clone();
49                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
50                })
51                .await
52                .context("reserve recovery for daemon resume task")??;
53                blocking({
54                    let session_id = session_id.clone();
55                    move || {
56                        crate::database::reassign_resumable_session_workspace(
57                            &session_id,
58                            &rebind_workspace_id,
59                        )
60                    }
61                })
62                .await?;
63                let restore_request = request.clone();
64                let mut controller = tokio::task::spawn_blocking(move || {
65                    session_move::load_controller_for_resume(&restore_request)
66                })
67                .await
68                .context("load controller for daemon resume task")??;
69                let executor = DaemonStageReportingExecutor::new(
70                    CancellableProcessExecutor::new(cancelled),
71                    state.clone(),
72                    session_id.clone(),
73                );
74                let materialized = controller
75                    .resume_session_controlled_with_repository_preflight(
76                        &session_id,
77                        &request.profile_id,
78                        &request.target_template_id,
79                        SessionResumeOptions {
80                            additional_mounts: request.additional_mounts,
81                            resource_allocation: request.resource_allocation,
82                            discard_queue: request.discard_queue,
83                        },
84                        request.repository_preflight,
85                        &executor,
86                    )
87                    .await?;
88                // The projection stays where it was written. A viewer reads
89                // it from the store; shipping it back through the daemon
90                // reply put a whole transcript in one IPC frame.
91                let _ = materialized;
92                Ok(DaemonLifecycleResult::Done)
93            },
94        )?;
95        self.set_lifecycle_resume_destination(
96            &operation_session_id,
97            profile_id,
98            target_template_id,
99        );
100        let channel = result.clone();
101        let result = Self::wait_lifecycle_result(result).await;
102        self.remove_completed_lifecycle(&channel);
103        match result? {
104            DaemonLifecycleResult::Done => {}
105            DaemonLifecycleResult::Move(_) => unreachable!("resume cannot return a move outcome"),
106            DaemonLifecycleResult::DeferredCleanup => {
107                unreachable!("session resume cannot schedule target cleanup")
108            }
109        }
110        blocking(move || {
111            if let Some(mut operation) =
112                crate::database::load_move_operation(&operation_session_id)?
113                && !operation.queue_admission_started
114            {
115                operation.phase = mj_core::state::MovePhase::Cancelled;
116                operation.queue_admission_finished = true;
117                operation.updated_at = chrono::Utc::now().to_rfc3339();
118                operation.error = Some("Recovered through an explicit Resume operation".into());
119                crate::database::save_move_operation(&operation)?;
120            }
121            Ok(())
122        })
123        .await?;
124        Ok(())
125    }
126
127    pub(super) async fn discard_since_checkpoint(
128        self: &Arc<Self>,
129        session_id: String,
130        checkpoint: mj_core::state::CheckpointMetadata,
131    ) -> Result<()> {
132        let children = blocking({
133            let session_id = session_id.clone();
134            move || {
135                let controller = Controller::load()?;
136                Ok(active_child_session_ids(&controller.state, &session_id))
137            }
138        })
139        .await?;
140        for child_id in children {
141            Box::pin(self.suspend_session(child_id.clone()))
142                .await
143                .with_context(|| {
144                    format!("suspend sub-agent {child_id} before discarding parent changes")
145                })?;
146        }
147        let operation_session_id = session_id.clone();
148        let result = self
149            .run_lifecycle(
150                operation_session_id,
151                LifecycleKind::ForceStop,
152                move |state, session_id, cancelled| async move {
153                    blocking(move || {
154                        let mut controller = Controller::load()?;
155                        let executor = DaemonStageReportingExecutor::new(
156                            CancellableProcessExecutor::new(cancelled),
157                            state,
158                            session_id.clone(),
159                        );
160                        ensure!(
161                            controller
162                                .state
163                                .sessions
164                                .get(&session_id)
165                                .and_then(|s| s.checkpoint.as_ref())
166                                == Some(&checkpoint),
167                            "the recovery copy changed; review it before discarding changes"
168                        );
169                        let deferred = controller.force_stop(&session_id, &executor)?;
170                        Ok(if deferred {
171                            DaemonLifecycleResult::DeferredCleanup
172                        } else {
173                            DaemonLifecycleResult::Done
174                        })
175                    })
176                    .await
177                },
178            )
179            .await?;
180        let _ = result; // The lifecycle supervisor owns the cleanup handoff.
181        Ok(())
182    }
183
184    pub(super) async fn destroy_stopped_session(
185        self: &Arc<Self>,
186        session_id: String,
187        branch: BranchDisposition,
188    ) -> Result<()> {
189        self.tear_down_stopped_session(
190            session_id,
191            LifecycleKind::DestroyStopped,
192            branch,
193            CheckoutDisposition::Remove,
194        )
195        .await
196        .map(|_| ())
197    }
198
199    /// Discard the record of a session that ended as `Lost`: its managed
200    /// target is gone, it has no checkpoint to resume from, and the only
201    /// action it still offers is Destroy. Keeping it would leave a tombstone
202    /// the person has to find and delete by hand.
203    ///
204    /// The branch stays, and so does a checkout that holds uncommitted
205    /// changes: nothing here was confirmed by a person, so nothing here may
206    /// discard work. A teardown that fails leaves the row where it is, with
207    /// its cause, and says so.
208    pub(crate) async fn discard_lost_session(self: &Arc<Self>, session_id: String) {
209        let short = mj_core::state::short_id(&session_id).to_owned();
210        match self
211            .tear_down_stopped_session(
212                session_id.clone(),
213                LifecycleKind::DestroyStopped,
214                BranchDisposition::Keep,
215                CheckoutDisposition::KeepWhenDirty,
216            )
217            .await
218        {
219            Ok(retained) => {
220                let mut text = format!(
221                    "Session {short} was lost because its managed target no longer exists; its record was removed."
222                );
223                if let Some(path) = retained {
224                    text.push_str(&format!(
225                        " Its checkout has uncommitted changes, so it was kept at {}.",
226                        path.display()
227                    ));
228                }
229                self.push_notice(&session_id, text);
230            }
231            Err(error) => {
232                tracing::warn!(
233                    %session_id,
234                    error = format!("{error:#}"),
235                    "could not discard the record of a lost session"
236                );
237                self.push_notice(
238                    &session_id,
239                    format!(
240                        "Session {short} was lost, but its record could not be removed: {error:#}"
241                    ),
242                );
243            }
244        }
245    }
246
247    /// Archive every stopped session older than `older_than_days` whose
248    /// conversation SessionWiki holds. Answers with how many were archived.
249    ///
250    /// The whole job belongs in a background task: it runs a full index sync,
251    /// which walks every tool's store, and then one lifecycle per session.
252    pub(crate) async fn archive_aged_sessions(
253        self: &Arc<Self>,
254        older_than_days: u32,
255    ) -> Result<usize> {
256        self.wiki()
257            .sync_now(true)
258            .await
259            .context("sync SessionWiki before archiving stopped sessions")?;
260        let candidates = blocking(move || {
261            let controller = Controller::load()?;
262            Ok(crate::sessionwiki::sessions_ready_to_archive(
263                &controller.state.sessions,
264                &controller.state.subagents,
265                chrono::Utc::now(),
266                older_than_days,
267            ))
268        })
269        .await
270        .context("select the stopped sessions old enough to archive")?;
271        if candidates.is_empty() {
272            return Ok(0);
273        }
274        let indexed = blocking({
275            let candidates = candidates.clone();
276            move || crate::sessionwiki::indexed_with_messages(&candidates)
277        })
278        .await
279        .context("check the SessionWiki index before archiving")?;
280        let mut archived = 0;
281        for session_id in candidates {
282            if !indexed.contains(&session_id) {
283                tracing::warn!(
284                    %session_id,
285                    "SessionWiki holds no conversation for this stopped session; keeping it"
286                );
287                continue;
288            }
289            match self.archive_stopped_session(session_id.clone()).await {
290                Ok(()) => {
291                    archived += 1;
292                    tracing::info!(
293                        %session_id,
294                        older_than_days,
295                        "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
296                    );
297                }
298                Err(error) => tracing::warn!(
299                    %session_id,
300                    error = %format!("{error:#}"),
301                    "could not archive a stopped session"
302                ),
303            }
304        }
305        if archived > 0 {
306            // Flip the rows the job just emptied to archived.
307            self.wiki().request_sync(false);
308        }
309        Ok(archived)
310    }
311
312    /// Destroy a stopped session the way the archive job wants: the record,
313    /// the checkpoint, and the attachments go, and the session's git branch
314    /// goes only when another branch already contains all of its commits.
315    /// The conversation itself stays searchable, and restorable, through
316    /// SessionWiki.
317    async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
318        self.tear_down_stopped_session(
319            session_id,
320            LifecycleKind::ArchiveStopped,
321            BranchDisposition::DeleteIfMerged,
322            CheckoutDisposition::Remove,
323        )
324        .await
325        .map(|_| ())
326    }
327
328    /// Answers with the managed checkout the teardown kept, if it kept one.
329    async fn tear_down_stopped_session(
330        self: &Arc<Self>,
331        session_id: String,
332        kind: LifecycleKind,
333        branch: BranchDisposition,
334        checkout: CheckoutDisposition,
335    ) -> Result<Option<PathBuf>> {
336        let children = blocking({
337            let session_id = session_id.clone();
338            move || {
339                Ok(crate::database::list_subagents(&session_id)?
340                    .into_iter()
341                    .map(|child| child.child_session_id)
342                    .collect::<Vec<_>>())
343            }
344        })
345        .await?;
346        for child_id in children {
347            // A sub-agent borrows its parent's worker and never owns a managed
348            // worktree, so it has no branch of its own to keep.
349            Box::pin(self.force_destroy_session(child_id.clone(), BranchDisposition::Keep))
350                .await
351                .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
352        }
353        self.wait_for_deferred_cleanup(&session_id).await?;
354        let exists = blocking({
355            let session_id = session_id.clone();
356            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
357        })
358        .await?;
359        if !exists {
360            return Ok(None);
361        }
362        // The retained path is produced inside the lifecycle task, which can
363        // only answer with a `DaemonLifecycleResult`, so it comes back here.
364        let retained = Arc::new(Mutex::new(None));
365        self.run_lifecycle(session_id, kind, {
366            let retained = retained.clone();
367            move |state, session_id, cancelled| async move {
368                blocking(move || {
369                    let mut controller = Controller::load()?;
370                    let executor = DaemonStageReportingExecutor::new(
371                        CancellableProcessExecutor::new(cancelled),
372                        state,
373                        session_id.clone(),
374                    );
375                    let kept = controller.destroy_session_controlled_with_checkout(
376                        &session_id,
377                        &executor,
378                        branch,
379                        checkout,
380                    )?;
381                    *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
382                    Ok(DaemonLifecycleResult::Done)
383                })
384                .await
385            }
386        })
387        .await?;
388        let kept = retained
389            .lock()
390            .unwrap_or_else(PoisonError::into_inner)
391            .clone();
392        Ok(kept)
393    }
394
395    /// Cancel any in-flight lifecycle for `session_id` and wait for it to
396    /// finish.
397    ///
398    /// Force destruction is the escape hatch for a wedged operation, so it
399    /// takes over rather than queueing behind one — but only after the running
400    /// task has stopped, because a cancelled create or close re-persists its
401    /// record as it unwinds and would otherwise resurrect the row this
402    /// operation deletes. A lifecycle that ignores cancellation for longer
403    /// than [`FORCE_DESTROY_PREEMPT_TIMEOUT`] is reported instead of destroyed
404    /// under.
405    pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
406        let mut result = {
407            let lifecycle = self
408                .lifecycle
409                .lock()
410                .unwrap_or_else(PoisonError::into_inner);
411            let Some(active) = lifecycle.get(session_id) else {
412                return Ok(());
413            };
414            if !active.result.borrow().is_none() {
415                return Ok(());
416            }
417            active.request_cancel();
418            active.result.clone()
419        };
420        let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
421            loop {
422                if result.borrow().is_some() {
423                    return Ok(());
424                }
425                if result.changed().await.is_err() {
426                    return Err(());
427                }
428            }
429        })
430        .await;
431        match finished {
432            // The loop only returns once the watch holds a result or its
433            // sender died; distinguish those two, and the timeout separately.
434            Ok(Ok(())) => Ok(()),
435            Ok(Err(())) => bail!(
436                "daemon lifecycle operation stopped without a result for session {session_id}"
437            ),
438            Err(_) => bail!(
439                "session {session_id} still has an operation that did not stop after cancellation; try again"
440            ),
441        }
442    }
443
444    /// Permanently destroy a session from any state, cancelling whatever
445    /// lifecycle operation holds it first. Data loss is the caller's confirmed
446    /// decision; see [`Controller::force_destroy_session`]. The session's git
447    /// branch survives unless `branch` says to delete it.
448    pub async fn force_destroy_session(
449        self: &Arc<Self>,
450        session_id: String,
451        branch: BranchDisposition,
452    ) -> Result<()> {
453        let children = blocking({
454            let session_id = session_id.clone();
455            move || {
456                Ok(crate::database::list_subagents(&session_id)?
457                    .into_iter()
458                    .map(|child| child.child_session_id)
459                    .collect::<Vec<_>>())
460            }
461        })
462        .await?;
463        for child_id in children {
464            // Sub-agents borrow their parent's worker and own no branch.
465            Box::pin(self.force_destroy_session(child_id.clone(), BranchDisposition::Keep))
466                .await
467                .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
468        }
469        self.preempt_active_lifecycle(&session_id).await?;
470        let exists = blocking({
471            let session_id = session_id.clone();
472            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
473        })
474        .await?;
475        if !exists {
476            return Ok(());
477        }
478        self.run_lifecycle(
479            session_id,
480            LifecycleKind::ForceDestroy,
481            move |state, session_id, cancelled| async move {
482                let _recovery_reservation = tokio::task::spawn_blocking({
483                    let observer = state.recovery_observer.clone();
484                    let session_id = session_id.clone();
485                    let cancelled = cancelled.clone();
486                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
487                })
488                .await
489                .context("reserve recovery for daemon force-destroy task")??;
490                blocking({
491                    let session_id = session_id.clone();
492                    move || {
493                        let mut controller = Controller::load()?;
494                        let executor = DaemonStageReportingExecutor::new(
495                            CancellableProcessExecutor::new(cancelled),
496                            state,
497                            session_id.clone(),
498                        );
499                        controller.force_destroy_session(&session_id, &executor, branch)?;
500                        crate::controller::move_session::release_move_queue_hold(&session_id);
501                        Ok(DaemonLifecycleResult::Done)
502                    }
503                })
504                .await
505            },
506        )
507        .await?;
508        Ok(())
509    }
510
511    /// Force-delete a workspace: destroy every active session in it (see
512    /// [`RuntimeState::force_destroy_session`]), drop its detached drafts, and
513    /// remove the workspace row. Stopped histories stay globally resumable.
514    ///
515    /// In-flight resumes into the workspace still refuse the deletion because
516    /// they have not yet claimed a durable session workspace. A session that
517    /// fails to destroy stops the sequence with the remainder named, so the
518    /// operation can be retried without losing progress.
519    pub async fn force_delete_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
520        ensure!(
521            !self.workspace_has_active_resume(&workspace_id),
522            "workspace has a session resume in progress"
523        );
524        let sessions = blocking({
525            let workspace_id = workspace_id.clone();
526            move || {
527                let controller = Controller::load()?;
528                Ok(active_sessions_for_force_destruction(
529                    &controller,
530                    &workspace_id,
531                ))
532            }
533        })
534        .await?;
535        for (index, session_id) in sessions.iter().enumerate() {
536            // Deleting a workspace removes Mjolnir's own copies, not the
537            // user's work: the branches stay in their source repositories.
538            if let Err(error) = self
539                .force_destroy_session(session_id.clone(), BranchDisposition::Keep)
540                .await
541            {
542                let remaining = sessions.len() - index - 1;
543                bail!(
544                    "force-destroying session {session_id} failed: {error:#}; \
545                     {remaining} session(s) in the workspace remain"
546                );
547            }
548        }
549        blocking({
550            let workspace_id = workspace_id.clone();
551            move || crate::database::force_delete_workspace(&workspace_id)
552        })
553        .await?;
554        refresh_runtime_workspaces(self).await?;
555        Ok(())
556    }
557}