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 relation = blocking(move || {
18            let mut controller = Controller::load()?;
19            controller.register_subagent(request)
20        })
21        .await?;
22        let session_id = relation.child_session_id.clone();
23        self.start_or_join_lifecycle_controlled(
24            session_id.clone(),
25            LifecycleKind::Create,
26            None,
27            Some(relation.request_key.clone()),
28            None,
29            move |state, session_id, cancelled| async move {
30                let mut controller = tokio::task::spawn_blocking(Controller::load)
31                    .await
32                    .context("load controller for sub-agent startup")??;
33                let executor = DaemonStageReportingExecutor::new(
34                    CancellableProcessExecutor::new(cancelled),
35                    state,
36                    session_id.clone(),
37                );
38                controller
39                    .provision_subagent_session_controlled(&session_id, &executor)
40                    .await?;
41                Ok(DaemonLifecycleResult::Done)
42            },
43        )?;
44        self.reload_controller().await?;
45        Ok(relation)
46    }
47
48    pub async fn start_create_session_controlled(
49        self: &Arc<Self>,
50        request: CreateSessionRequest,
51        control: CreateSessionControl,
52        publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
53    ) -> Result<RegisteredSession> {
54        self.start_create_session_inner(request, control, Some(publication))
55            .await
56    }
57
58    pub(super) async fn start_create_session_inner(
59        self: &Arc<Self>,
60        request: CreateSessionRequest,
61        control: CreateSessionControl,
62        publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
63    ) -> Result<RegisteredSession> {
64        let path_cancelled = control.cancelled.clone();
65        let registered = blocking(move || {
66            let mut controller = Controller::load()?;
67            let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
68                .with_deadline(Duration::from_secs(30));
69            let project_directory = request
70                .project_directory
71                .as_deref()
72                .map(|path| {
73                    controller.resolve_project_directory(
74                        &request.target_template_id,
75                        path,
76                        &path_executor,
77                    )
78                })
79                .transpose()?;
80            let session_id = controller.register_session_with_resources(
81                &request.profile_id,
82                &request.bundle_id,
83                &request.target_template_id,
84                request.title,
85                SessionLaunchOptions {
86                    create_managed_worktree: request.create_managed_worktree,
87                    mjolnir_subagents: request.mjolnir_subagents,
88                    initial_prompt: request.initial_prompt,
89                    workspace_id: request.workspace_id,
90                    additional_mounts: request.additional_mounts,
91                    resource_allocation: request.resource_allocation,
92                    project_directory,
93                    session_title_override: request.session_title_override,
94                },
95            )?;
96            let session = controller
97                .state
98                .sessions
99                .get(&session_id)
100                .expect("newly registered session exists")
101                .clone();
102            let remembered_container_size = controller
103                .config
104                .targets
105                .get(&request.target_template_id)
106                .and_then(mj_core::config::container_size_host)
107                .and_then(|host| {
108                    controller
109                        .state
110                        .container_sizes
111                        .get(host)
112                        .copied()
113                        .map(|size| (host.to_owned(), size))
114                });
115            Ok(RegisteredSession {
116                session,
117                remembered_container_size,
118            })
119        })
120        .await?;
121        let session_id = registered.session.id.clone();
122        self.start_or_join_lifecycle_controlled(
123            session_id,
124            LifecycleKind::Create,
125            None,
126            None,
127            Some(control.clone()),
128            move |state, session_id, cancelled| async move {
129                let mut controller = tokio::task::spawn_blocking(Controller::load)
130                    .await
131                    .context("load controller for daemon create task")??;
132                let publication_error = if let Some(publication) = publication {
133                    let published = tokio::select! {
134                        result = publication => result.context("session publication owner stopped")
135                            .and_then(|result| result.map_err(anyhow::Error::msg)),
136                        () = async {
137                            while !cancelled.load(Ordering::Acquire) {
138                                tokio::time::sleep(Duration::from_millis(25)).await;
139                            }
140                        } => Err(anyhow!("session creation cancelled before publication")),
141                    };
142                    published.err()
143                } else {
144                    None
145                };
146                if publication_error.is_some() {
147                    control.request_cancel();
148                }
149                let executor = DaemonStageReportingExecutor::new(
150                    CancellableProcessExecutor::new(cancelled),
151                    state,
152                    session_id.clone(),
153                );
154                let provision = controller
155                    .provision_session_controlled_with_commit(&session_id, &executor, || {
156                        ensure!(
157                            control.grant_commit(),
158                            "session creation cancelled before commit"
159                        );
160                        Ok(())
161                    })
162                    .await;
163                if let Some(error) = publication_error {
164                    return match provision {
165                        Ok(()) => Err(error),
166                        Err(rollback) => {
167                            Err(error.context(format!("discard unpublished session: {rollback:#}")))
168                        }
169                    };
170                }
171                provision?;
172                Ok(DaemonLifecycleResult::Done)
173            },
174        )?;
175        self.reload_controller().await?;
176        Ok(registered)
177    }
178
179    pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
180        let result = {
181            let lifecycle = self
182                .lifecycle
183                .lock()
184                .unwrap_or_else(PoisonError::into_inner);
185            let active = lifecycle
186                .get(session_id)
187                .with_context(|| format!("no create operation exists for session {session_id}"))?;
188            ensure!(
189                active.kind == LifecycleKind::Create,
190                "session {session_id} is no longer being created"
191            );
192            active.result.clone()
193        };
194        let channel = result.clone();
195        let outcome = Self::wait_lifecycle_result(result).await;
196        self.remove_completed_lifecycle(&channel);
197        match outcome? {
198            DaemonLifecycleResult::Done => Ok(()),
199            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
200            DaemonLifecycleResult::DeferredCleanup => {
201                unreachable!("session creation cannot schedule target cleanup")
202            }
203        }
204    }
205}