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                    mjolnir_subagents: request.mjolnir_subagents,
90                    initial_prompt: request.initial_prompt,
91                    workspace_id: request.workspace_id,
92                    additional_mounts: request.additional_mounts,
93                    resource_allocation: request.resource_allocation,
94                    project_directory,
95                    session_title_override: request.session_title_override,
96                },
97            )?;
98            let session = controller
99                .state
100                .sessions
101                .get(&session_id)
102                .expect("newly registered session exists")
103                .clone();
104            let remembered_container_size = controller
105                .config
106                .targets
107                .get(&request.target_template_id)
108                .and_then(mj_core::config::container_size_host)
109                .and_then(|host| {
110                    controller
111                        .state
112                        .container_sizes
113                        .get(host)
114                        .copied()
115                        .map(|size| (host.to_owned(), size))
116                });
117            Ok(RegisteredSession {
118                session,
119                remembered_container_size,
120            })
121        })
122        .await?;
123        let session_id = registered.session.id.clone();
124        self.start_or_join_lifecycle_controlled(
125            session_id,
126            LifecycleKind::Create,
127            None,
128            None,
129            Some(control.clone()),
130            move |state, session_id, cancelled| async move {
131                let mut controller = tokio::task::spawn_blocking(Controller::load)
132                    .await
133                    .context("load controller for daemon create task")??;
134                let publication_error = if let Some(publication) = publication {
135                    let published = tokio::select! {
136                        result = publication => result.context("session publication owner stopped")
137                            .and_then(|result| result.map_err(anyhow::Error::msg)),
138                        () = async {
139                            while !cancelled.load(Ordering::Acquire) {
140                                tokio::time::sleep(Duration::from_millis(25)).await;
141                            }
142                        } => Err(anyhow!("session creation cancelled before publication")),
143                    };
144                    published.err()
145                } else {
146                    None
147                };
148                if publication_error.is_some() {
149                    control.request_cancel();
150                }
151                let executor = DaemonStageReportingExecutor::new(
152                    CancellableProcessExecutor::new(cancelled),
153                    state,
154                    session_id.clone(),
155                );
156                let provision = controller
157                    .provision_session_controlled_with_commit(&session_id, &executor, || {
158                        ensure!(
159                            control.grant_commit(),
160                            "session creation cancelled before commit"
161                        );
162                        Ok(())
163                    })
164                    .await;
165                if let Some(error) = publication_error {
166                    return match provision {
167                        Ok(()) => Err(error),
168                        Err(rollback) => {
169                            Err(error.context(format!("discard unpublished session: {rollback:#}")))
170                        }
171                    };
172                }
173                provision?;
174                Ok(DaemonLifecycleResult::Done)
175            },
176        )?;
177        self.reload_controller().await?;
178        Ok(registered)
179    }
180
181    pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
182        let result = {
183            let lifecycle = self
184                .lifecycle
185                .lock()
186                .unwrap_or_else(PoisonError::into_inner);
187            let active = lifecycle
188                .get(session_id)
189                .with_context(|| format!("no create operation exists for session {session_id}"))?;
190            ensure!(
191                active.kind == LifecycleKind::Create,
192                "session {session_id} is no longer being created"
193            );
194            active.result.clone()
195        };
196        let channel = result.clone();
197        let outcome = Self::wait_lifecycle_result(result).await;
198        self.remove_completed_lifecycle(&channel);
199        match outcome? {
200            DaemonLifecycleResult::Done => Ok(()),
201            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
202            DaemonLifecycleResult::DeferredCleanup => {
203                unreachable!("session creation cannot schedule target cleanup")
204            }
205        }
206    }
207}