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        let owned_clone = blocking({
211            let session_id = session_id.clone();
212            move || {
213                let controller = Controller::load()?;
214                let Some(checkout) = controller
215                    .state
216                    .sessions
217                    .get(&session_id)
218                    .and_then(|session| session.managed_worktree.as_ref())
219                    .filter(|checkout| checkout.kind == mj_core::state::ManagedCheckoutKind::Clone)
220                else {
221                    return Ok(None);
222                };
223                if crate::controller::path_exists_on_managed_target(
224                    &crate::targets::ProcessExecutor,
225                    &checkout.target,
226                    &checkout.worktree_root,
227                )? {
228                    Ok(Some(checkout.worktree_root.clone()))
229                } else {
230                    Ok(None)
231                }
232            }
233        })
234        .await;
235        match owned_clone {
236            Ok(Some(path)) => {
237                self.push_notice(&session_id, format!(
238                    "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
239                    path.display(),
240                ));
241                return;
242            }
243            Err(error) => {
244                tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
245                return;
246            }
247            Ok(None) => {}
248        }
249        match self
250            .tear_down_stopped_session(
251                session_id.clone(),
252                LifecycleKind::DestroyStopped,
253                BranchDisposition::Keep,
254                CheckoutDisposition::KeepWhenDirty,
255            )
256            .await
257        {
258            Ok(retained) => {
259                let mut text = format!(
260                    "Session {short} was lost because its managed target no longer exists; its record was removed."
261                );
262                if let Some(path) = retained {
263                    text.push_str(&format!(
264                        " Its checkout has uncommitted changes, so it was kept at {}.",
265                        path.display()
266                    ));
267                }
268                self.push_notice(&session_id, text);
269            }
270            Err(error) => {
271                tracing::warn!(
272                    %session_id,
273                    error = format!("{error:#}"),
274                    "could not discard the record of a lost session"
275                );
276                self.push_notice(
277                    &session_id,
278                    format!(
279                        "Session {short} was lost, but its record could not be removed: {error:#}"
280                    ),
281                );
282            }
283        }
284    }
285
286    /// Archive every stopped session older than `older_than_days` whose
287    /// conversation SessionWiki holds. Answers with how many were archived.
288    ///
289    /// The whole job belongs in a background task: it runs a full index sync,
290    /// which walks every tool's store, and then one lifecycle per session.
291    pub(crate) async fn archive_aged_sessions(
292        self: &Arc<Self>,
293        older_than_days: u32,
294    ) -> Result<usize> {
295        self.wiki()
296            .sync_now(true)
297            .await
298            .context("sync SessionWiki before archiving stopped sessions")?;
299        // Recheck a small bounded batch of older clone checkpoints. A remote
300        // that was offline during suspension may now prove its saved refs.
301        let refreshable = blocking(move || {
302            let controller = Controller::load()?;
303            let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
304            let mut sessions = controller
305                .state
306                .sessions
307                .values()
308                .filter(|session| session.state == SessionState::Stopped)
309                .filter(|session| {
310                    chrono::DateTime::parse_from_rfc3339(&session.updated_at)
311                        .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
312                })
313                .filter(|session| {
314                    session.managed_worktree.as_ref().is_some_and(|owned| {
315                        owned.kind == mj_core::state::ManagedCheckoutKind::Clone
316                    })
317                })
318                .filter(|session| {
319                    session.publication.as_ref().is_none_or(|evidence| {
320                        evidence.state != mj_core::state::PublicationState::Published
321                            && !evidence.dirty
322                            && !evidence.stashed
323                    })
324                })
325                .cloned()
326                .collect::<Vec<_>>();
327            sessions.sort_by(|a, b| {
328                a.publication
329                    .as_ref()
330                    .map(|e| &e.checked_at)
331                    .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
332            });
333            sessions.truncate(4);
334            Ok(sessions)
335        })
336        .await?;
337        let mut checks = tokio::task::JoinSet::new();
338        for session in refreshable {
339            checks.spawn_blocking(move || {
340                let assessment =
341                    crate::controller::publication::refresh_stopped_clone_publication(&session);
342                (session.id, assessment)
343            });
344        }
345        let mut refreshed = false;
346        while let Some(done) = checks.join_next().await {
347            let (id, assessment) = done.context("publication refresh task failed")?;
348            if let Some(assessment) = assessment {
349                refreshed |= blocking(move || {
350                    crate::database::set_publication_assessment_if_current(&id, &assessment)
351                })
352                .await?;
353            }
354        }
355        if refreshed {
356            self.reload_controller().await?;
357            self.publish_revision();
358        }
359        let candidates = blocking(move || {
360            let controller = Controller::load()?;
361            Ok(crate::sessionwiki::sessions_ready_to_archive(
362                &controller.state.sessions,
363                &controller.state.subagents,
364                chrono::Utc::now(),
365                older_than_days,
366            ))
367        })
368        .await
369        .context("select the stopped sessions old enough to archive")?;
370        if candidates.is_empty() {
371            return Ok(0);
372        }
373        let indexed = blocking({
374            let candidates = candidates.clone();
375            move || crate::sessionwiki::indexed_with_messages(&candidates)
376        })
377        .await
378        .context("check the SessionWiki index before archiving")?;
379        let mut archived = 0;
380        for session_id in candidates {
381            if !indexed.contains(&session_id) {
382                tracing::warn!(
383                    %session_id,
384                    "SessionWiki holds no conversation for this stopped session; keeping it"
385                );
386                continue;
387            }
388            match self.archive_stopped_session(session_id.clone()).await {
389                Ok(()) => {
390                    archived += 1;
391                    tracing::info!(
392                        %session_id,
393                        older_than_days,
394                        "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
395                    );
396                }
397                Err(error) => tracing::warn!(
398                    %session_id,
399                    error = %format!("{error:#}"),
400                    "could not archive a stopped session"
401                ),
402            }
403        }
404        if archived > 0 {
405            // Flip the rows the job just emptied to archived.
406            self.wiki().request_sync(false);
407        }
408        Ok(archived)
409    }
410
411    /// Destroy a stopped session the way the archive job wants: the record,
412    /// the checkpoint, and the attachments go, and the session's git branch
413    /// goes only when another branch already contains all of its commits.
414    /// The conversation itself stays searchable, and restorable, through
415    /// SessionWiki.
416    async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
417        self.tear_down_stopped_session(
418            session_id,
419            LifecycleKind::ArchiveStopped,
420            BranchDisposition::DeleteIfMerged,
421            CheckoutDisposition::Remove,
422        )
423        .await
424        .map(|_| ())
425    }
426
427    /// Answers with the managed checkout the teardown kept, if it kept one.
428    async fn tear_down_stopped_session(
429        self: &Arc<Self>,
430        session_id: String,
431        kind: LifecycleKind,
432        branch: BranchDisposition,
433        checkout: CheckoutDisposition,
434    ) -> Result<Option<PathBuf>> {
435        // This indexes the sub-agents too, so they are destroyed below
436        // without indexing each one again.
437        self.index_before_destroy(&session_id).await;
438        let children = blocking({
439            let session_id = session_id.clone();
440            move || {
441                Ok(crate::database::list_subagents(&session_id)?
442                    .into_iter()
443                    .map(|child| child.child_session_id)
444                    .collect::<Vec<_>>())
445            }
446        })
447        .await?;
448        for child_id in children {
449            // A sub-agent borrows its parent's worker and never owns a managed
450            // worktree, so it has no branch of its own to keep.
451            Box::pin(self.force_destroy_indexed_session(child_id.clone(), BranchDisposition::Keep))
452                .await
453                .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
454        }
455        self.wait_for_deferred_cleanup(&session_id).await?;
456        let exists = blocking({
457            let session_id = session_id.clone();
458            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
459        })
460        .await?;
461        if !exists {
462            return Ok(None);
463        }
464        // The retained path is produced inside the lifecycle task, which can
465        // only answer with a `DaemonLifecycleResult`, so it comes back here.
466        let retained = Arc::new(Mutex::new(None));
467        self.run_lifecycle(session_id, kind, {
468            let retained = retained.clone();
469            move |state, session_id, cancelled| async move {
470                blocking(move || {
471                    let mut controller = Controller::load()?;
472                    let executor = DaemonStageReportingExecutor::new(
473                        CancellableProcessExecutor::new(cancelled),
474                        state,
475                        session_id.clone(),
476                    );
477                    let kept = controller.destroy_session_controlled_with_checkout(
478                        &session_id,
479                        &executor,
480                        branch,
481                        checkout,
482                    )?;
483                    *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
484                    Ok(DaemonLifecycleResult::Done)
485                })
486                .await
487            }
488        })
489        .await?;
490        let kept = retained
491            .lock()
492            .unwrap_or_else(PoisonError::into_inner)
493            .clone();
494        Ok(kept)
495    }
496
497    /// Put a session about to be destroyed, and its sub-agents, into
498    /// SessionWiki while their records and stored conversations still exist,
499    /// so `mj sessions --session <id>` still finds them afterwards (R2-11).
500    ///
501    /// Waits at most [`crate::sessionwiki::DESTROY_SYNC_WAIT`] for a sync
502    /// pass, then indexes the sessions on their own. A destroy is never
503    /// refused for this: when the index cannot take the sessions, the log
504    /// says why and the destroy goes ahead.
505    async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
506        use crate::sessionwiki::IndexedBeforeDestroy;
507        let outcome = self
508            .wiki()
509            .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
510            .await;
511        match outcome {
512            IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
513                %session_id,
514                reason,
515                "destroying a session without indexing it in SessionWiki"
516            ),
517            IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
518                %session_id,
519                %reason,
520                "could not index a session in SessionWiki before destroying it"
521            ),
522            outcome => tracing::debug!(
523                %session_id,
524                ?outcome,
525                "indexed a session in SessionWiki before destroying it"
526            ),
527        }
528    }
529
530    /// Cancel any in-flight lifecycle for `session_id` and wait for it to
531    /// finish.
532    ///
533    /// Force destruction is the escape hatch for a wedged operation, so it
534    /// takes over rather than queueing behind one — but only after the running
535    /// task has stopped, because a cancelled create or close re-persists its
536    /// record as it unwinds and would otherwise resurrect the row this
537    /// operation deletes. A lifecycle that ignores cancellation for longer
538    /// than [`FORCE_DESTROY_PREEMPT_TIMEOUT`] is reported instead of destroyed
539    /// under.
540    pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
541        let mut result = {
542            let lifecycle = self
543                .lifecycle
544                .lock()
545                .unwrap_or_else(PoisonError::into_inner);
546            let Some(active) = lifecycle.get(session_id) else {
547                return Ok(());
548            };
549            if !active.result.borrow().is_none() {
550                return Ok(());
551            }
552            active.request_cancel();
553            active.result.clone()
554        };
555        let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
556            loop {
557                if result.borrow().is_some() {
558                    return Ok(());
559                }
560                if result.changed().await.is_err() {
561                    return Err(());
562                }
563            }
564        })
565        .await;
566        match finished {
567            // The loop only returns once the watch holds a result or its
568            // sender died; distinguish those two, and the timeout separately.
569            Ok(Ok(())) => Ok(()),
570            Ok(Err(())) => bail!(
571                "daemon lifecycle operation stopped without a result for session {session_id}"
572            ),
573            Err(_) => bail!(
574                "session {session_id} still has an operation that did not stop after cancellation; try again"
575            ),
576        }
577    }
578
579    /// Permanently destroy a session from any state, cancelling whatever
580    /// lifecycle operation holds it first. Data loss is the caller's confirmed
581    /// decision; see [`Controller::force_destroy_session`]. The session's git
582    /// branch survives unless `branch` says to delete it.
583    pub async fn force_destroy_session(
584        self: &Arc<Self>,
585        session_id: String,
586        branch: BranchDisposition,
587    ) -> Result<()> {
588        self.index_before_destroy(&session_id).await;
589        self.force_destroy_indexed_session(session_id, branch).await
590    }
591
592    /// [`Self::force_destroy_session`] once the session and its sub-agents
593    /// have been indexed.
594    async fn force_destroy_indexed_session(
595        self: &Arc<Self>,
596        session_id: String,
597        branch: BranchDisposition,
598    ) -> Result<()> {
599        let children = blocking({
600            let session_id = session_id.clone();
601            move || {
602                Ok(crate::database::list_subagents(&session_id)?
603                    .into_iter()
604                    .map(|child| child.child_session_id)
605                    .collect::<Vec<_>>())
606            }
607        })
608        .await?;
609        for child_id in children {
610            // Sub-agents borrow their parent's worker and own no branch.
611            Box::pin(self.force_destroy_indexed_session(child_id.clone(), BranchDisposition::Keep))
612                .await
613                .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
614        }
615        self.preempt_active_lifecycle(&session_id).await?;
616        let exists = blocking({
617            let session_id = session_id.clone();
618            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
619        })
620        .await?;
621        if !exists {
622            return Ok(());
623        }
624        self.run_lifecycle(
625            session_id,
626            LifecycleKind::ForceDestroy,
627            move |state, session_id, cancelled| async move {
628                let _recovery_reservation = tokio::task::spawn_blocking({
629                    let observer = state.recovery_observer.clone();
630                    let session_id = session_id.clone();
631                    let cancelled = cancelled.clone();
632                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
633                })
634                .await
635                .context("reserve recovery for daemon force-destroy task")??;
636                blocking({
637                    let session_id = session_id.clone();
638                    move || {
639                        let mut controller = Controller::load()?;
640                        let executor = DaemonStageReportingExecutor::new(
641                            CancellableProcessExecutor::new(cancelled),
642                            state,
643                            session_id.clone(),
644                        );
645                        controller.force_destroy_session(&session_id, &executor, branch)?;
646                        crate::controller::move_session::release_move_queue_hold(&session_id);
647                        Ok(DaemonLifecycleResult::Done)
648                    }
649                })
650                .await
651            },
652        )
653        .await?;
654        Ok(())
655    }
656
657    /// Force-delete a workspace: destroy every active session in it (see
658    /// [`RuntimeState::force_destroy_session`]), drop its detached drafts, and
659    /// remove the workspace row. Stopped histories stay globally resumable.
660    ///
661    /// In-flight resumes into the workspace still refuse the deletion because
662    /// they have not yet claimed a durable session workspace. A session that
663    /// fails to destroy stops the sequence with the remainder named, so the
664    /// operation can be retried without losing progress.
665    pub async fn force_delete_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
666        ensure!(
667            !self.workspace_has_active_resume(&workspace_id),
668            "workspace has a session resume in progress"
669        );
670        let sessions = blocking({
671            let workspace_id = workspace_id.clone();
672            move || {
673                let controller = Controller::load()?;
674                Ok(active_sessions_for_force_destruction(
675                    &controller,
676                    &workspace_id,
677                ))
678            }
679        })
680        .await?;
681        for (index, session_id) in sessions.iter().enumerate() {
682            // Deleting a workspace removes Mjolnir's own copies, not the
683            // user's work: the branches stay in their source repositories.
684            if let Err(error) = self
685                .force_destroy_session(session_id.clone(), BranchDisposition::Keep)
686                .await
687            {
688                let remaining = sessions.len() - index - 1;
689                // The verb agrees with the count as well as the noun.
690                let remaining = if remaining == 1 {
691                    "1 session in the workspace remains".to_owned()
692                } else {
693                    format!("{remaining} sessions in the workspace remain")
694                };
695                bail!("force-destroying session {session_id} failed: {error:#}; {remaining}");
696            }
697        }
698        blocking({
699            let workspace_id = workspace_id.clone();
700            move || crate::database::force_delete_workspace(&workspace_id)
701        })
702        .await?;
703        refresh_runtime_workspaces(self).await?;
704        Ok(())
705    }
706}