mj_controller/daemon/
close_workspace.rs1use 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
18fn 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 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 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 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}