Skip to main content

mj_controller/daemon/
create.rs

1use super::*;
2
3/// How long creating a parent's report root on its target may take.
4const 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    /// Register and start a child worker on its parent's existing target.
16    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            // The parent's report root is made before the child exists, so the
38            // child's first prompt can name its own directory under it. A
39            // request run again after a restart finds the same root.
40            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                // The lifecycle owner now excludes competing resume/close work.
66                // Recovery may already have settled this registration before
67                // admission; replay never reinstalls an existing worker.
68                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        // Discovery is restartable preparation, not admitted lifecycle work.
133        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            // Every surface creates sessions through here, so the dashboard,
172            // the phone, `mj new`, and `mj acp` all leave a default pair
173            // behind for the next caller that names none. A preference that
174            // cannot be written does not undo a session that was created.
175            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}