Skip to main content

mj_controller/daemon/
create.rs

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