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