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        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                // The lifecycle owner now excludes competing resume/close work.
67                // Recovery may already have settled this registration before
68                // admission; replay never reinstalls an existing worker.
69                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        // Discovery is restartable preparation, not admitted lifecycle work.
165        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            // Every surface creates sessions through here, so the dashboard,
204            // the phone, `mj new`, and `mj acp` all leave a default pair
205            // behind for the next caller that names none. A preference that
206            // cannot be written does not undo a session that was created.
207            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}