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