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 = blocking(move || {
22            let mut controller = Controller::load()?;
23            let mut request = request;
24            // The parent's report root is made before the child exists, so the
25            // child's first prompt can name its own directory under it. A
26            // request run again after a restart finds the same root.
27            let executor = CancellableProcessExecutor::with_timeout(REPORT_ROOT_TIMEOUT);
28            request.report_root = Some(
29                controller
30                    .prepare_subagent_report_root(&request.parent_session_id, &executor)
31                    .context("create the sub-agent report directory")?,
32            );
33            controller.register_subagent(request)
34        })
35        .await?;
36        let session_id = relation.child_session_id.clone();
37        self.start_or_join_lifecycle_controlled(
38            session_id.clone(),
39            LifecycleKind::Create,
40            None,
41            Some(relation.request_key.clone()),
42            None,
43            move |state, session_id, cancelled| async move {
44                let mut controller = tokio::task::spawn_blocking(Controller::load)
45                    .await
46                    .context("load controller for sub-agent startup")??;
47                let executor = DaemonStageReportingExecutor::new(
48                    CancellableProcessExecutor::new(cancelled),
49                    state,
50                    session_id.clone(),
51                );
52                controller
53                    .provision_subagent_session_controlled(&session_id, &executor)
54                    .await?;
55                Ok(DaemonLifecycleResult::Done)
56            },
57        )?;
58        self.reload_controller().await?;
59        Ok(relation)
60    }
61
62    pub async fn start_create_session_controlled(
63        self: &Arc<Self>,
64        request: CreateSessionRequest,
65        control: CreateSessionControl,
66        publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
67    ) -> Result<RegisteredSession> {
68        self.start_create_session_inner(request, control, Some(publication))
69            .await
70    }
71
72    pub(super) async fn start_create_session_inner(
73        self: &Arc<Self>,
74        request: CreateSessionRequest,
75        control: CreateSessionControl,
76        publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
77    ) -> Result<RegisteredSession> {
78        let _upgrade_work = crate::upgrade::activity("session admission")?;
79        let path_cancelled = control.cancelled.clone();
80        let registered = blocking(move || {
81            let mut controller = Controller::load()?;
82            let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
83                .with_deadline(Duration::from_secs(30));
84            let project_directory = request
85                .project_directory
86                .as_deref()
87                .map(|path| {
88                    controller.resolve_project_directory(
89                        &request.target_template_id,
90                        path,
91                        &path_executor,
92                    )
93                })
94                .transpose()?;
95            let session_id = controller.register_session_with_resources(
96                &request.profile_id,
97                &request.bundle_id,
98                &request.target_template_id,
99                request.title,
100                SessionLaunchOptions {
101                    create_managed_worktree: request.create_managed_worktree,
102                    launch_base: request.launch_base,
103                    launch_branch: request.launch_branch,
104                    checkout: request.checkout,
105                    expected_runtime_identity: request.expected_runtime_identity,
106                    mjolnir_subagents: request.mjolnir_subagents,
107                    initial_prompt: request.initial_prompt,
108                    workspace_id: request.workspace_id,
109                    additional_mounts: request.additional_mounts,
110                    resource_allocation: request.resource_allocation,
111                    project_directory,
112                    session_title_override: request.session_title_override,
113                },
114            )?;
115            // Every surface creates sessions through here, so the dashboard,
116            // the phone, `mj new`, and `mj acp` all leave a default pair
117            // behind for the next caller that names none. A preference that
118            // cannot be written does not undo a session that was created.
119            if let Err(error) = mj_core::go::GoPreferences::remember_first_pair(
120                &mj_core::go::GoPreferences::path(),
121                &request.profile_id,
122                &request.target_template_id,
123            ) {
124                tracing::warn!(%error, "could not save the default profile and target");
125            }
126            let session = controller
127                .state
128                .sessions
129                .get(&session_id)
130                .expect("newly registered session exists")
131                .clone();
132            let remembered_container_size = controller
133                .config
134                .targets
135                .get(&request.target_template_id)
136                .and_then(mj_core::config::container_size_host)
137                .and_then(|host| {
138                    controller
139                        .state
140                        .container_sizes
141                        .get(host)
142                        .copied()
143                        .map(|size| (host.to_owned(), size))
144                });
145            Ok(RegisteredSession {
146                session,
147                remembered_container_size,
148            })
149        })
150        .await?;
151        let session_id = registered.session.id.clone();
152        self.start_or_join_lifecycle_controlled(
153            session_id,
154            LifecycleKind::Create,
155            None,
156            None,
157            Some(control.clone()),
158            move |state, session_id, cancelled| async move {
159                let mut controller = tokio::task::spawn_blocking(Controller::load)
160                    .await
161                    .context("load controller for daemon create task")??;
162                let publication_error = if let Some(publication) = publication {
163                    let published = tokio::select! {
164                        result = publication => result.context("session publication owner stopped")
165                            .and_then(|result| result.map_err(anyhow::Error::msg)),
166                        () = async {
167                            while !cancelled.load(Ordering::Acquire) {
168                                tokio::time::sleep(Duration::from_millis(25)).await;
169                            }
170                        } => Err(anyhow!("session creation cancelled before publication")),
171                    };
172                    published.err()
173                } else {
174                    None
175                };
176                if publication_error.is_some() {
177                    control.request_cancel();
178                }
179                let executor = DaemonStageReportingExecutor::new(
180                    CancellableProcessExecutor::new(cancelled),
181                    state,
182                    session_id.clone(),
183                );
184                let provision = controller
185                    .provision_session_controlled_with_commit(&session_id, &executor, || {
186                        ensure!(
187                            control.grant_commit(),
188                            "session creation cancelled before commit"
189                        );
190                        Ok(())
191                    })
192                    .await;
193                if let Some(error) = publication_error {
194                    return match provision {
195                        Ok(()) => Err(error),
196                        Err(rollback) => {
197                            Err(error.context(format!("discard unpublished session: {rollback:#}")))
198                        }
199                    };
200                }
201                provision?;
202                Ok(DaemonLifecycleResult::Done)
203            },
204        )?;
205        self.reload_controller().await?;
206        Ok(registered)
207    }
208
209    pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
210        let result = {
211            let lifecycle = self
212                .lifecycle
213                .lock()
214                .unwrap_or_else(PoisonError::into_inner);
215            let active = lifecycle
216                .get(session_id)
217                .with_context(|| format!("no create operation exists for session {session_id}"))?;
218            ensure!(
219                active.kind == LifecycleKind::Create,
220                "session {session_id} is no longer being created"
221            );
222            active.result.clone()
223        };
224        let channel = result.clone();
225        let outcome = Self::wait_lifecycle_result(result).await;
226        self.remove_completed_lifecycle(&channel);
227        match outcome? {
228            DaemonLifecycleResult::Done => Ok(()),
229            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
230            DaemonLifecycleResult::DeferredCleanup => {
231                unreachable!("session creation cannot schedule target cleanup")
232            }
233        }
234    }
235}