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 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 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 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 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}