Skip to main content

mj_controller/daemon/
close_workspace.rs

1use super::*;
2
3struct WorkspaceCloseGuard {
4    state: Arc<RuntimeState>,
5    workspace_id: String,
6}
7
8impl Drop for WorkspaceCloseGuard {
9    fn drop(&mut self) {
10        self.state
11            .workspace_closes
12            .lock()
13            .unwrap_or_else(PoisonError::into_inner)
14            .remove(&self.workspace_id);
15    }
16}
17
18/// Only roots are submitted: normal close already stops their active children.
19fn close_roots(state: &mj_core::state::State, workspace_id: &str) -> Vec<String> {
20    let active: BTreeSet<_> = state
21        .sessions
22        .values()
23        .filter(|session| session.workspace_id == workspace_id && session.state.is_active())
24        .map(|session| session.id.clone())
25        .collect();
26    active
27        .iter()
28        .filter(|id| {
29            !state.subagents.values().any(|child| {
30                &child.child_session_id == *id && active.contains(&child.parent_session_id)
31            })
32        })
33        .cloned()
34        .collect()
35}
36
37impl RuntimeState {
38    pub(super) fn workspace_resume_gate(&self, workspace_id: &str) -> Arc<tokio::sync::RwLock<()>> {
39        self.workspace_resume_admission
40            .lock()
41            .unwrap_or_else(PoisonError::into_inner)
42            .entry(workspace_id.to_owned())
43            .or_default()
44            .clone()
45    }
46
47    pub(super) fn cancel_workspace_close(&self, workspace_id: &str) -> Result<()> {
48        let closes = self
49            .workspace_closes
50            .lock()
51            .unwrap_or_else(PoisonError::into_inner);
52        let cancelled = closes
53            .get(workspace_id)
54            .context("workspace close is no longer running")?;
55        cancelled.store(true, Ordering::Release);
56        Ok(())
57    }
58
59    pub async fn close_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
60        let _upgrade_work = crate::upgrade::activity("workspace close")?;
61        let cancelled = Arc::new(AtomicBool::new(false));
62        {
63            let mut closes = self
64                .workspace_closes
65                .lock()
66                .unwrap_or_else(PoisonError::into_inner);
67            ensure!(
68                !closes.contains_key(&workspace_id),
69                "workspace is already closing"
70            );
71            closes.insert(workspace_id.clone(), cancelled.clone());
72        }
73        let _guard = WorkspaceCloseGuard {
74            state: self.clone(),
75            workspace_id: workspace_id.clone(),
76        };
77        ensure!(
78            !self.workspace_has_active_resume(&workspace_id),
79            "workspace has a session resume in progress"
80        );
81        let (roots, sessions) = blocking({
82            let workspace_id = workspace_id.clone();
83            move || {
84                let controller = Controller::load()?;
85                let sessions = controller
86                    .state
87                    .sessions
88                    .values()
89                    .filter(|session| {
90                        session.workspace_id == workspace_id && session.state.is_active()
91                    })
92                    .map(|session| session.id.clone())
93                    .collect::<Vec<_>>();
94                Ok((close_roots(&controller.state, &workspace_id), sessions))
95            }
96        })
97        .await?;
98        let mut jobs = tokio::task::JoinSet::new();
99        for session_id in roots {
100            let state = self.clone();
101            let cancelled = cancelled.clone();
102            jobs.spawn(async move {
103                ensure!(
104                    !cancelled.load(Ordering::Acquire),
105                    "workspace close cancelled"
106                );
107                state
108                    .suspend_session(session_id.clone())
109                    .await
110                    .with_context(|| format!("stop session {session_id}"))
111            });
112        }
113        let mut failures = Vec::new();
114        let mut cancellation_poll = tokio::time::interval(Duration::from_millis(100));
115        while !jobs.is_empty() {
116            tokio::select! {
117                result = jobs.join_next() => match result {
118                    Some(Ok(Ok(()))) => {},
119                    Some(Ok(Err(error))) => failures.push(format!("{error:#}")),
120                    Some(Err(error)) => failures.push(format!("workspace close task failed: {error}")),
121                    None => break,
122                },
123                _ = cancellation_poll.tick() => {
124                    if cancelled.load(Ordering::Acquire) {
125                        for session_id in &sessions {
126                            // A close past checkpoint commit must finish its teardown.
127                            if let Err(error) = self.cancel_lifecycle(session_id) {
128                                tracing::debug!(%session_id, %error, "workspace close cancellation cannot interrupt this session");
129                            }
130                        }
131                    }
132                }
133            }
134        }
135        ensure!(
136            failures.is_empty(),
137            "Workspace retained; some sessions may already be stopped. Retry to close remaining sessions: {}",
138            failures.join("; ")
139        );
140        ensure!(
141            !cancelled.load(Ordering::Acquire),
142            "Workspace deletion cancelled; workspace and drafts retained. Completed suspensions were not undone."
143        );
144        // A resumed session may still have its old workspace id in storage.
145        // Exclude admission until the deletion transaction has committed.
146        let _admission = self
147            .workspace_resume_gate(&workspace_id)
148            .try_write_owned()
149            .context("a session resume is in progress; workspace retained, retry closing it")?;
150        ensure!(
151            !self.workspace_has_active_resume(&workspace_id),
152            "workspace has a session resume in progress"
153        );
154        blocking({
155            let workspace_id = workspace_id.clone();
156            let cancelled = cancelled.clone();
157            move || {
158                ensure!(
159                    !cancelled.load(Ordering::Acquire),
160                    "workspace close cancelled; workspace retained"
161                );
162                // The transaction rejects new active sessions; insertion triggers
163                // reject sessions registered after their workspace disappears.
164                crate::database::close_workspace(&workspace_id)
165            }
166        })
167        .await?;
168        refresh_runtime_workspaces(self).await?;
169        Ok(())
170    }
171}
172
173#[cfg(test)]
174mod tests {
175    use super::*;
176
177    #[test]
178    fn close_roots_do_not_submit_children_twice_or_touch_other_workspaces() {
179        let mut state = mj_core::state::State::default();
180        for (id, workspace, status) in [
181            ("parent", "a", SessionState::Running),
182            ("child", "a", SessionState::Running),
183            ("independent", "a", SessionState::Running),
184            ("history", "a", SessionState::Stopped),
185            ("other", "b", SessionState::Running),
186        ] {
187            let session = super::super::tests::runtime_test_session(id, workspace, status);
188            state.sessions.insert(id.into(), session);
189        }
190        let child = super::super::tests::runtime_test_subagent("child", "parent");
191        state.subagents.insert("child".into(), child);
192        assert_eq!(close_roots(&state, "a"), ["independent", "parent"]);
193        state.sessions.get_mut("parent").unwrap().state = SessionState::Stopped;
194        assert_eq!(close_roots(&state, "a"), ["child", "independent"]);
195    }
196}