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