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