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 cancelled = Arc::new(AtomicBool::new(false));
61        {
62            let mut closes = self
63                .workspace_closes
64                .lock()
65                .unwrap_or_else(PoisonError::into_inner);
66            ensure!(
67                !closes.contains_key(&workspace_id),
68                "workspace is already closing"
69            );
70            closes.insert(workspace_id.clone(), cancelled.clone());
71        }
72        let _guard = WorkspaceCloseGuard {
73            state: self.clone(),
74            workspace_id: workspace_id.clone(),
75        };
76        ensure!(
77            !self.workspace_has_active_resume(&workspace_id),
78            "workspace has a session resume in progress"
79        );
80        let (roots, sessions) = blocking({
81            let workspace_id = workspace_id.clone();
82            move || {
83                let controller = Controller::load()?;
84                let sessions = controller
85                    .state
86                    .sessions
87                    .values()
88                    .filter(|session| {
89                        session.workspace_id == workspace_id && session.state.is_active()
90                    })
91                    .map(|session| session.id.clone())
92                    .collect::<Vec<_>>();
93                Ok((close_roots(&controller.state, &workspace_id), sessions))
94            }
95        })
96        .await?;
97        let mut jobs = tokio::task::JoinSet::new();
98        for session_id in roots {
99            let state = self.clone();
100            let cancelled = cancelled.clone();
101            jobs.spawn(async move {
102                ensure!(
103                    !cancelled.load(Ordering::Acquire),
104                    "workspace close cancelled"
105                );
106                state
107                    .suspend_session(session_id.clone())
108                    .await
109                    .with_context(|| format!("stop session {session_id}"))
110            });
111        }
112        let mut failures = Vec::new();
113        let mut cancellation_poll = tokio::time::interval(Duration::from_millis(100));
114        while !jobs.is_empty() {
115            tokio::select! {
116                result = jobs.join_next() => match result {
117                    Some(Ok(Ok(()))) => {},
118                    Some(Ok(Err(error))) => failures.push(format!("{error:#}")),
119                    Some(Err(error)) => failures.push(format!("workspace close task failed: {error}")),
120                    None => break,
121                },
122                _ = cancellation_poll.tick() => {
123                    if cancelled.load(Ordering::Acquire) {
124                        for session_id in &sessions {
125                            // A close past checkpoint commit must finish its teardown.
126                            if let Err(error) = self.cancel_lifecycle(session_id) {
127                                tracing::debug!(%session_id, %error, "workspace close cancellation cannot interrupt this session");
128                            }
129                        }
130                    }
131                }
132            }
133        }
134        ensure!(
135            failures.is_empty(),
136            "Workspace retained; some sessions may already be stopped. Retry to close remaining sessions: {}",
137            failures.join("; ")
138        );
139        ensure!(
140            !cancelled.load(Ordering::Acquire),
141            "Workspace close cancelled; workspace and drafts retained. Completed stops were not undone."
142        );
143        // A resumed session may still have its old workspace id in storage.
144        // Exclude admission until the deletion transaction has committed.
145        let _admission = self
146            .workspace_resume_gate(&workspace_id)
147            .try_write_owned()
148            .context("a session resume is in progress; workspace retained, retry closing it")?;
149        ensure!(
150            !self.workspace_has_active_resume(&workspace_id),
151            "workspace has a session resume in progress"
152        );
153        blocking({
154            let workspace_id = workspace_id.clone();
155            let cancelled = cancelled.clone();
156            move || {
157                ensure!(
158                    !cancelled.load(Ordering::Acquire),
159                    "workspace close cancelled; workspace retained"
160                );
161                // The transaction rejects new active sessions; insertion triggers
162                // reject sessions registered after their workspace disappears.
163                crate::database::close_workspace(&workspace_id)
164            }
165        })
166        .await?;
167        refresh_runtime_workspaces(self).await?;
168        Ok(())
169    }
170}
171
172#[cfg(test)]
173mod tests {
174    use super::*;
175
176    #[test]
177    fn close_roots_do_not_submit_children_twice_or_touch_other_workspaces() {
178        let mut state = mj_core::state::State::default();
179        for (id, workspace, status) in [
180            ("parent", "a", SessionState::Running),
181            ("child", "a", SessionState::Running),
182            ("independent", "a", SessionState::Running),
183            ("history", "a", SessionState::Stopped),
184            ("other", "b", SessionState::Running),
185        ] {
186            let session = super::super::tests::runtime_test_session(id, workspace, status);
187            state.sessions.insert(id.into(), session);
188        }
189        let child = super::super::tests::runtime_test_subagent("child", "parent");
190        state.subagents.insert("child".into(), child);
191        assert_eq!(close_roots(&state, "a"), ["independent", "parent"]);
192        state.sessions.get_mut("parent").unwrap().state = SessionState::Stopped;
193        assert_eq!(close_roots(&state, "a"), ["child", "independent"]);
194    }
195}