1use super::*;
2
3const REPORT_ROOT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
5
6impl RuntimeState {
7 pub(super) async fn start_create_session(
8 self: &Arc<Self>,
9 request: CreateSessionRequest,
10 ) -> Result<RegisteredSession> {
11 self.start_create_session_inner(request, CreateSessionControl::default(), None)
12 .await
13 }
14
15 pub async fn start_subagent_session(
17 self: &Arc<Self>,
18 request: crate::controller::RegisterSubagentRequest,
19 ) -> Result<mj_core::subagent::SubagentRecord> {
20 let _upgrade_work = crate::upgrade::activity("subagent admission")?;
21 let (relation, needs_provisioning) = blocking(move || {
22 let mut controller = Controller::load()?;
23 if let Some(existing) = crate::database::lookup_subagent_request(
24 &request.parent_session_id,
25 &request.request_key,
26 )? {
27 let needs_provisioning = controller
28 .state
29 .sessions
30 .get(&existing.child_session_id)
31 .is_some_and(|record| {
32 record.state == mj_core::state::SessionState::Provisioning
33 });
34 return Ok((existing, needs_provisioning));
35 }
36 let mut request = request;
37 let executor = CancellableProcessExecutor::with_timeout(REPORT_ROOT_TIMEOUT);
41 request.report_root = Some(
42 controller
43 .prepare_subagent_report_root(&request.parent_session_id, &executor)
44 .context("create the sub-agent report directory")?,
45 );
46 controller
47 .register_subagent(request)
48 .map(|relation| (relation, true))
49 })
50 .await?;
51 if !needs_provisioning {
52 return Ok(relation);
53 }
54 let session_id = relation.child_session_id.clone();
55 let parent_session_id = relation.parent_session_id.clone();
56 let started = self.start_or_join_lifecycle_controlled(
57 session_id.clone(),
58 LifecycleKind::Create,
59 None,
60 Some(relation.request_key.clone()),
61 None,
62 move |state, session_id, cancelled| async move {
63 let mut controller = tokio::task::spawn_blocking(Controller::load)
64 .await
65 .context("load controller for sub-agent startup")??;
66 if controller
70 .state
71 .sessions
72 .get(&session_id)
73 .is_none_or(|record| record.state != mj_core::state::SessionState::Provisioning)
74 {
75 return Ok(DaemonLifecycleResult::Done);
76 }
77 let executor = DaemonStageReportingExecutor::new(
78 CancellableProcessExecutor::new(cancelled),
79 state.clone(),
80 session_id.clone(),
81 );
82 if let Err(error) = controller
83 .provision_subagent_session_controlled(&session_id, &executor)
84 .await
85 {
86 drop(controller);
87 state.reload_controller().await?;
88 if let Err(prompt_error) = state
89 .ensure_parent_wait_prompt(&parent_session_id)
90 .await
91 {
92 tracing::warn!(
93 child_session_id = %session_id,
94 parent_session_id,
95 error = format!("{prompt_error:#}"),
96 "could not reconcile the parent's sub-agent wait prompt after provisioning failure"
97 );
98 }
99 return Err(error);
100 }
101 Ok(DaemonLifecycleResult::Done)
102 },
103 );
104 if let Err(error) = started {
105 self.reload_controller().await?;
106 if let Err(prompt_error) = self
107 .ensure_parent_wait_prompt(&relation.parent_session_id)
108 .await
109 {
110 tracing::warn!(
111 child_session_id = %relation.child_session_id,
112 parent_session_id = %relation.parent_session_id,
113 error = format!("{prompt_error:#}"),
114 "could not reconcile the parent's sub-agent wait prompt after provisioning failure"
115 );
116 }
117 return Err(error);
118 }
119 self.reload_controller().await?;
120 Ok(relation)
121 }
122
123 pub async fn start_create_session_controlled(
124 self: &Arc<Self>,
125 request: CreateSessionRequest,
126 control: CreateSessionControl,
127 publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
128 ) -> Result<RegisteredSession> {
129 self.start_create_session_inner(request, control, Some(publication))
130 .await
131 }
132
133 pub(super) async fn start_create_session_inner(
134 self: &Arc<Self>,
135 request: CreateSessionRequest,
136 control: CreateSessionControl,
137 publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
138 ) -> Result<RegisteredSession> {
139 let mut request = request;
140 let parent = request.profile_id.clone();
141 let supplied = request.subagents.clone();
142 let (config, policy) = blocking(move || {
143 let controller = Controller::load()?;
144 let profile = controller
145 .config
146 .enabled_profile(&parent)
147 .context("parent profile unavailable")?;
148 let policy = supplied.unwrap_or_else(|| profile.subagents.clone());
149 Ok((controller.config, policy))
150 })
151 .await?;
152 crate::controller::profile_config::validate_session_subagent_policy(
153 &config,
154 &request.profile_id,
155 &policy,
156 )
157 .await?;
158 if !request.bundle_id.is_empty() {
159 crate::controller::config_only_controller(config.clone())
160 .validate_github_bundle_installations(&request.bundle_id)
161 .await
162 .map_err(crate::controller::GithubBundleSelectionError::into_anyhow)?;
163 }
164 let _upgrade_work = crate::upgrade::activity("session admission")?;
166 request.subagents = Some(policy);
167 let path_cancelled = control.cancelled.clone();
168 let registered = blocking(move || {
169 let mut controller = Controller::load()?;
170 let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
171 .with_deadline(Duration::from_secs(30));
172 let project_directory = request
173 .project_directory
174 .as_deref()
175 .map(|path| {
176 controller.resolve_project_directory(
177 &request.target_template_id,
178 path,
179 &path_executor,
180 )
181 })
182 .transpose()?;
183 let session_id = controller.register_session_with_resources(
184 &request.profile_id,
185 &request.bundle_id,
186 &request.target_template_id,
187 request.title,
188 SessionLaunchOptions {
189 create_managed_worktree: request.create_managed_worktree,
190 at: request.at,
191 branch: request.branch,
192 base: request.base,
193 subagents: request.subagents,
194 review: request.review,
195 initial_prompt: request.initial_prompt,
196 workspace_id: request.workspace_id,
197 additional_mounts: request.additional_mounts,
198 resource_allocation: request.resource_allocation,
199 project_directory,
200 session_title_override: request.session_title_override,
201 },
202 )?;
203 if let Err(error) = mj_core::go::GoPreferences::remember_first_pair(
208 &mj_core::go::GoPreferences::path(),
209 &request.profile_id,
210 &request.target_template_id,
211 ) {
212 tracing::warn!(%error, "could not save the default profile and target");
213 }
214 let session = controller
215 .state
216 .sessions
217 .get(&session_id)
218 .expect("newly registered session exists")
219 .clone();
220 let remembered_container_size = controller
221 .config
222 .targets
223 .get(&request.target_template_id)
224 .and_then(mj_core::config::container_size_host)
225 .and_then(|host| {
226 controller
227 .state
228 .container_sizes
229 .get(host)
230 .copied()
231 .map(|size| (host.to_owned(), size))
232 });
233 Ok(RegisteredSession {
234 session,
235 remembered_container_size,
236 })
237 })
238 .await?;
239 let session_id = registered.session.id.clone();
240 self.start_or_join_lifecycle_controlled(
241 session_id,
242 LifecycleKind::Create,
243 None,
244 None,
245 Some(control.clone()),
246 move |state, session_id, cancelled| async move {
247 let mut controller = tokio::task::spawn_blocking(Controller::load)
248 .await
249 .context("load controller for daemon create task")??;
250 let publication_error = if let Some(publication) = publication {
251 let published = tokio::select! {
252 result = publication => result.context("session publication owner stopped")
253 .and_then(|result| result.map_err(anyhow::Error::msg)),
254 () = async {
255 while !cancelled.load(Ordering::Acquire) {
256 tokio::time::sleep(Duration::from_millis(25)).await;
257 }
258 } => Err(anyhow!("session creation cancelled before publication")),
259 };
260 published.err()
261 } else {
262 None
263 };
264 if publication_error.is_some() {
265 control.request_cancel();
266 }
267 let executor = DaemonStageReportingExecutor::new(
268 CancellableProcessExecutor::new(cancelled),
269 state,
270 session_id.clone(),
271 );
272 let provision = controller
273 .provision_session_controlled_with_commit(&session_id, &executor, || {
274 ensure!(
275 control.grant_commit(),
276 "session creation cancelled before commit"
277 );
278 Ok(())
279 })
280 .await;
281 if let Some(error) = publication_error {
282 return match provision {
283 Ok(()) => Err(error),
284 Err(rollback) => {
285 Err(error.context(format!("discard unpublished session: {rollback:#}")))
286 }
287 };
288 }
289 provision?;
290 Ok(DaemonLifecycleResult::Done)
291 },
292 )?;
293 self.reload_controller().await?;
294 Ok(registered)
295 }
296
297 pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
298 let result = {
299 let lifecycle_owner = self.owner();
300 let lifecycle = &lifecycle_owner.lifecycle;
301 let active = lifecycle
302 .get(session_id)
303 .with_context(|| format!("no create operation exists for session {session_id}"))?;
304 ensure!(
305 active.kind == LifecycleKind::Create,
306 "session {session_id} is no longer being created"
307 );
308 active.result.clone()
309 };
310 let channel = result.clone();
311 let outcome = Self::wait_lifecycle_result(result).await;
312 self.remove_completed_lifecycle(&channel);
313 match outcome? {
314 DaemonLifecycleResult::Done => Ok(()),
315 DaemonLifecycleResult::Move(_) | DaemonLifecycleResult::Park(_) => {
316 unreachable!("cleanup cannot return a move outcome")
317 }
318 DaemonLifecycleResult::DeferredCleanup => {
319 unreachable!("session creation cannot schedule target cleanup")
320 }
321 }
322 }
323}
324
325#[cfg(test)]
326mod delegation_replay_tests {
327 use crate::controller::test_support::{IsolatedTest, test_name};
328
329 #[tokio::test]
330 async fn multi_model_creation_is_refused_before_registration_or_provisioning() {
331 const CHILD: &str = "MJ_TEST_MULTI_MODEL_CREATION_REFUSAL";
332 if std::env::var_os(CHILD).is_none() {
333 let root = tempfile::tempdir().unwrap();
334 IsolatedTest::new(test_name(
335 module_path!(),
336 "multi_model_creation_is_refused_before_registration_or_provisioning",
337 ))
338 .env(CHILD, "1")
339 .env("MJ_INSTANCE", "multi-model-creation-refusal")
340 .isolated_store(root.path())
341 .run();
342 return;
343 }
344 let _writer = crate::database::install_isolated_test_writer();
345 let mut config = mj_core::config::Config::default();
346 config.profiles.insert(
347 "codex".into(),
348 mj_core::config::HarnessProfile {
349 enabled: true,
350 kind: mj_core::config::HarnessKind::Codex,
351 home: mj_core::config::data_dir().join("missing-profile-home"),
352 environment: Default::default(),
353 context_window_bytes: None,
354 guardian_review_model: None,
355 subagents: mj_core::subagent::SubagentPolicy::Native,
356 },
357 );
358 config.save().unwrap();
359 let workspace = crate::database::create_workspace("legacy delegation").unwrap();
360 let mut legacy = crate::daemon::tests::runtime_test_session(
361 "legacy-parent",
362 &workspace.id,
363 mj_core::state::SessionState::Running,
364 );
365 legacy.subagents = Some(mj_core::subagent::SubagentPolicy::AllModels);
366 crate::database::save_session(&legacy).unwrap();
367 let runtime = crate::daemon::tests::test_runtime_state();
368 let request = mj_client::daemon::CreateSessionRequest {
369 create_managed_worktree: None,
370 at: None,
371 branch: None,
372 base: None,
373 subagents: Some(mj_core::subagent::SubagentPolicy::AllModels),
374 review: None,
375 initial_prompt: None,
376 workspace_id: workspace.id,
377 profile_id: "codex".into(),
378 bundle_id: "missing-bundle".into(),
379 project_directory: None,
380 target_template_id: "missing-target".into(),
381 additional_mounts: Vec::new(),
382 resource_allocation: None,
383 title: "must not register".into(),
384 session_title_override: None,
385 };
386 let error = runtime
387 .start_create_session(request)
388 .await
389 .expect_err("a refusal");
390 let refusal = mj_core::refusal::Refusal::of(&error).expect("a request refusal");
391 assert!(refusal.message().contains("no longer available"));
392 assert!(runtime.active_lifecycles().is_empty());
393 let state = crate::database::load_state().unwrap();
394 assert_eq!(state.sessions.len(), 1);
395 assert_eq!(state.sessions[&legacy.id], legacy);
396 }
397
398 #[tokio::test]
399 async fn replayed_spawn_does_not_reprovision_an_existing_child() {
400 const CHILD: &str = "MJ_TEST_SPAWN_REPLAY";
401 if std::env::var_os(CHILD).is_none() {
402 let root = tempfile::tempdir().unwrap();
403 IsolatedTest::new(test_name(
404 module_path!(),
405 "replayed_spawn_does_not_reprovision_an_existing_child",
406 ))
407 .env(CHILD, "1")
408 .env("MJ_INSTANCE", "concurrency-sweep-spawn")
409 .isolated_store(root.path())
410 .run();
411 return;
412 }
413 let _writer = crate::database::install_isolated_test_writer();
414 let workspace = crate::database::create_workspace("spawn replay").unwrap();
415 let parent = crate::daemon::tests::runtime_test_session(
416 "parent",
417 &workspace.id,
418 mj_core::state::SessionState::Running,
419 );
420 crate::database::save_session(&parent).unwrap();
421 for (index, status) in [
422 mj_core::state::SessionState::Running,
423 mj_core::state::SessionState::Parked,
424 mj_core::state::SessionState::Stopped,
425 ]
426 .into_iter()
427 .enumerate()
428 {
429 let child_id = format!("child-{index}");
430 let child =
431 crate::daemon::tests::runtime_test_session(&child_id, &workspace.id, status);
432 let relation = crate::daemon::tests::runtime_test_subagent(&child_id, "parent");
433 crate::database::save_subagent_session(&child, &relation).unwrap();
434 let runtime = crate::daemon::tests::test_runtime_state();
435 let replayed = runtime
436 .start_subagent_session(crate::controller::RegisterSubagentRequest {
437 parent_session_id: "parent".into(),
438 task_name: "must reuse original".into(),
439 profile_id: "missing-profile".into(),
440 model: None,
441 effort: None,
442 working_directory: Default::default(),
443 initial_prompt: "must not resend".into(),
444 request_key: relation.request_key.clone(),
445 report_root: None,
446 })
447 .await
448 .unwrap();
449 assert_eq!(replayed, relation);
450 assert!(runtime.active_lifecycles().is_empty());
451 assert_eq!(
452 crate::database::load_session_record(&child_id)
453 .unwrap()
454 .unwrap()
455 .state,
456 status
457 );
458 }
459 }
460}